火花输出到卡夫卡一次

bfo*_*vdr 5 scala apache-kafka apache-spark

我想输出火花和火花流到卡夫卡一次.但正如文档所说的那样 "输出操作(如foreachRDD)至少有一次语义,也就是说,如果发生工作失败,转换后的数据可能会被多次写入外部实体."
要进行事务更新,spark建议使用批处理时间(在foreachRDD中可用)和RDD的分区索引来创建标识符.该标识符唯一地标识流应用程序中的blob数据.代码如下:

dstream.foreachRDD { (rdd, time) =>
  rdd.foreachPartition { partitionIterator =>
    val partitionId = TaskContext.get.partitionId()
    val **uniqueId** = generateUniqueId(time.milliseconds, partitionId)
    // use this uniqueId to transactionally commit the data in  partitionIterator
  }
}
Run Code Online (Sandbox Code Playgroud)

但是,如何使用UNIQUEID卡夫卡,使事务提交.

谢谢

cod*_*ure 1

Kixer 的高级软件工程师 Cody Koeninger 在 Spark 峰会上讨论了 Kafka 的一次性解决方案。本质上,该解决方案涉及通过同时提交来存储偏移量和数据。

在 2016 年的 Confluence 聚会上,工程师们在向工程师提到过一次主题时,引用了 Cody 关于该主题的演讲。Cloudera 在http://blog.cloudera.com/blog/2015/03/exactly-once-spark-streaming-from-apache-kafka/上发表了他的演讲。Cody 的论文位于http://koeninger.github.io/kafka-exactly-once/#1,他的 github(针对此主题)位于https://github.com/koeninger/kafka-exactly-once。网上还可以找到他的演讲视频。

Kafka 的更高版本引入了 Kafka Streams 来处理没有 Spark 的一次性场景,但该主题只值得一个脚注,因为问题的框架是与 Spark 一起使用。