【问题标题】:How to map kafka topic names and respective records in spark streaming如何在火花流中映射kafka主题名称和各自的记录
【发布时间】:2016-09-29 08:20:52
【问题描述】:

我正在流式传输 kafka 主题,如下所示;

JavaPairInputDStream<String, String> directKafkaStream = 
    KafkaUtils.createDirectStream(jssc,
                                  String.class, 
                                  String.class,
                                  StringDecoder.class,
                                  StringDecoder.class,
                                  kafkaParams, 
                                  topicSet);

directKafkaStream.print();   

一个主题的输出如下所示:

(null,"04/15/2015","18:44:14")
(null,"04/15/2015","18:44:15")
(null,"04/15/2015","18:44:16")
(null,"04/15/2015","18:44:17")  

如何映射主题名称和记录。
例如:主题是“callData”,它应该类似于下面等等

(callData,"04/15/2015","18:44:14")
(callData,"04/15/2015","18:44:15")
(callData,"04/15/2015","18:44:16")
(callData,"04/15/2015","18:44:17")  

【问题讨论】:

  • "map topic name and records"是什么意思?
  • 刚刚更新了问题..

标签: java apache-spark apache-kafka


【解决方案1】:

如何映射主题名称和记录?

为了提取分区信息,you'll need to use the overload which accepts a Function receiving MessageAndMetadata&lt;K, V&gt; 并返回您希望转换为的类型。

看起来像这样:

Map<TopicAndPartition, Long> map = new HashMap<>();
map.put(new TopicAndPartition("topicname", 0), 1L);

JavaInputDStream<Map.Entry> stream = KafkaUtils.createDirectStream(
        javaContext,
        String.class,
        String.class,
        StringDecoder.class,
        StringDecoder.class,
        Map.Entry.class, // <--- This is the record return type from the transformation.
        kafkaParams,
        map,
        messageAndMetadata -> 
            new AbstractMap.SimpleEntry<>(messageAndMetadata.topic(),
                                          messageAndMetadata.message()));

请注意,我使用 Map.Entry 作为 Java 替代 Scala 中的 Tuple2。您可以提供自己的类,该类也具有PartitionMessage 属性,并将 用于转换。请注意,kafka 输入流的类型现在是 JavaInputDStream&lt;Map.Entry&gt;,因为这是转换返回的内容。

【讨论】:

  • 16/05/31 10:44:20 错误 kafka.DirectKafkaInputDStream: ArrayBuffer(org.apache.spark.SparkException: 找不到 Set()) 的领导者偏移量
  • 这里map的值应该是什么。上面写着Map。我应该在这里放什么
  • @Alka 应该分别是你的队列、分区和偏移的名称。我已经编辑了一个示例。
  • 太棒了..谢谢...我需要更多信息..因为我是 spark 新手..从哪里获得有关 spark.which api 在什么地方使用的这些信息。跨度>
  • @Alka 很多关于跟踪和错误。您可以使用Spark API DocumentationSpark Documentation。如果您还有其他问题,我建议您打开一个新问题。
猜你喜欢
  • 2019-10-11
  • 2015-10-13
  • 1970-01-01
  • 2017-04-27
  • 2020-03-23
  • 2016-08-21
  • 1970-01-01
  • 1970-01-01
  • 2017-10-22
相关资源
最近更新 更多