Uh oh!
There was an error while loading. Please reload this page.
[SPARK-12177][Streaming][Kafka] Update KafkaDStreams to new Kafka 0.10 Consumer API - #11863
[SPARK-12177][Streaming][Kafka] Update KafkaDStreams to new Kafka 0.10 Consumer API#11863koeninger wants to merge 42 commits into
Conversation
… beta consumer, modify getPreferredLocations to choose a consistent executor per topicpartition
… for new consumers is finished
…g, but dont handle recalculating the same RDD efficiently
…mer for dynamic topics, listener, etc
…consuming messages on driver
SparkQA
commented
Mar 21, 2016
Test build #53680 has finished for PR 11863 at commit
|
SparkQA
commented
Mar 21, 2016
Test build #53682 has finished for PR 11863 at commit
|
SparkQA
commented
Mar 21, 2016
Test build #53686 has started for PR 11863 at commit |
…on attempts to serialize ConsumerRecord
SparkQA
commented
Apr 7, 2016
Test build #55252 has finished for PR 11863 at commit
|
| <groupId>org.apache.spark</groupId> | ||
| <artifactId>spark-streaming-kafka-beta-assembly_2.11</artifactId> | ||
| <packaging>jar</packaging> | ||
| <name>Spark Project External Kafka Assembly</name> |
There was a problem hiding this comment.
I think it may be a good idea to update this, so the two kafka assemblies can be differentiated in the build.
| @Experimental | ||
| object KafkaUtils extends Logging { | ||
| /** | ||
| * Scala constructor for a batch-oriented interface for consuming from Kafka. |
There was a problem hiding this comment.
Please add :: Experimental :: at the beginning of comments if you add the @Experimental tag.
zsxwing
commented
Jun 29, 2016
Finished my round of reviewing. Some some nits and one question about |
koeninger
commented
Jun 29, 2016
@zsxwing Thanks for the fixes |
| * configuration parameters</a>. | ||
| * Requires "bootstrap.servers" to be set with Kafka broker(s), | ||
| * NOT zookeeper servers, specified in host1:port1,host2:port2 form. | ||
| * @param driverConsumer zero-argument function for you to construct a Kafka Consumer, |
Overall, this is looking good. Two high level points.
|
koeninger
commented
Jun 30, 2016
You do need CanCommitOffsets because DirectKafkaInputDstream is now
|
SparkQA
commented
Jun 30, 2016
Test build #61495 has finished for PR 11863 at commit
|
tdas
commented
Jun 30, 2016
Aah, right. My bad. In that case, there arent major issues as far as i can see, let me merge this, and test how the docs look like. I am pretty sure its going to cause trouble with two KafkaUtils. And in that case I will handle the package renaming. |
tdas
commented
Jun 30, 2016
Well.. after the tests pass. |
koeninger
commented
Jun 30, 2016
I'll do the scaladoc fix and the package rename. I think the package rename is fine even if it did work with docs, just to disambiguate things. Will start a separate ticket for documentation updates. |
tdas
commented
Jun 30, 2016
sounds good. thanks! |
…10 version number, to disambiguate from the older connector
SparkQA
commented
Jun 30, 2016
Test build #61506 has finished for PR 11863 at commit
|
SparkQA
commented
Jun 30, 2016
Test build #3151 has finished for PR 11863 at commit
|
SparkQA
commented
Jun 30, 2016
Test build #3150 has finished for PR 11863 at commit
|
SparkQA
commented
Jun 30, 2016
Test build #61513 has finished for PR 11863 at commit
|
tdas
commented
Jun 30, 2016
LGTM. Merging this to master and 2.0. Thank you very much @koeninger for this awesome effort. :) |
…0 Consumer API ## What changes were proposed in this pull request? New Kafka consumer api for the released 0.10 version of Kafka ## How was this patch tested? Unit tests, manual tests Author: cody koeninger <cody@koeninger.org> Closes#11863 from koeninger/kafka-0.9. (cherry picked from commit dedbcee) Signed-off-by: Tathagata Das <tathagata.das1565@gmail.com>
| // make sure constructors can be called from java | ||
| final ConsumerStrategy<String, String> sub0 = | ||
| Subscribe.<String, String>apply(topics, kafkaParams, offsets); |
There was a problem hiding this comment.
This is seems to break in scala 2.10 and not scala 2.11. This is very weird.
Merging this PR broke 2.10 builds - https://amplab.cs.berkeley.edu/jenkins/view/Spark%20QA%20Compile/job/spark-master-compile-sbt-scala-2.10/1947/console
[error] /home/jenkins/workspace/spark-master-compile-sbt-scala-2.10/external/kafka-0-10/src/test/java/org/apache/spark/streaming/kafka010/JavaConsumerStrategySuite.java:54: error: incompatible types: Collection<String> cannot be converted to Iterable<String>
[error] Subscribe.<String, String>apply(topics, kafkaParams, offsets);
[error] ^
[error] /home/jenkins/workspace/spark-master-compile-sbt-scala-2.10/external/kafka-0-10/src/test/java/org/apache/spark/streaming/kafka010/JavaConsumerStrategySuite.java:69: error: incompatible types: Collection<TopicPartition> cannot be converted to Iterable<TopicPartition>
[error] Assign.<String, String>apply(parts, kafkaParams, offsets);
[error] ^
There was a problem hiding this comment.
We should figure out a way to fix scala 2.10. I don't think we need to revert this though since 2.10 is no longer the default build and it does not fail PRs.
There was a problem hiding this comment.
Okay found the issue. In scala 2.10, if companion object of a case class has explicitly defined apply(), then the implicit apply method is not generated. In scala 2.11 it is generated.
I remember now, this type of stuff is why we avoid using case classes in the public API. Do you mind if I convert these to simple classes??
There was a problem hiding this comment.
I refactored the API to avoid case classes and minimize publicly visible classes - #13996
tdas
commented
Jun 30, 2016
I played around with the API and I found a few issues
I have opened a new PR to address them - please take a look - #13996 |
hey,I have an question about the setting "enable.auto.commit",could it be changed ????Because I wanna save the offsets information to zookeeper cluster. |
koeninger
commented
Aug 3, 2017
via email
You won't get any reasonable semantics out of auto commit, because it will
commit on the driver without regard to what the executors have done. …On Aug 2, 2017 21:46, "Wallace Huang" ***@***.***> wrote:
hey,I have an question about the setting "auto.commit.enable", It could be
changed ????Because I wanna save the offsets information to zookeeper
cluster.
—
You are receiving this because you were mentioned.
Reply to this email directly, view it on GitHub
<#11863 (comment)>, or mute
the thread
<https://github.com/notifications/unsubscribe-auth/AAGAByVZ2QEf_o627Z7BKdRBuCtmza-Nks5sUTSTgaJpZM4H1Pg1>
.
|
BiyuHuang
commented
Aug 3, 2017
I'm wondering that why the setting "enable.auto.commit" existed, but it was set to false by default and I could't modify it . Anyway, how do I use it ? |
What changes were proposed in this pull request?
New Kafka consumer api for the released 0.10 version of Kafka
How was this patch tested?
Unit tests, manual tests