【问题标题】:kafka muliple topic seggregation in sparkspark中的kafka多主题聚合
【发布时间】:2016-05-30 21:09:37
【问题描述】:

我正在阅读来自 2 个不同主题的 2 个不同文件的 kafka 行。行示例:

例如: 文件 1:2015-04-15T18:44:14+01:00,192.168.11.42,%ASA-2-106007:
文件2:"04/15/2012","18:44:14",,"Start","Unknown","Unknown",,"192.168.63.128","444","2","7","192.168.63.128",,,,,,,,,,,,,,,,,

我可以从两个不同主题的火花中阅读。代码如下:

SparkConf sparkConfig = new SparkConf().setAppName("KafkaStreaming").setMaster("local[5]");
        JavaStreamingContext jsc = new JavaStreamingContext(sparkConfig,Durations.seconds(5));
        final HiveContext sqlContext = new HiveContext(jsc.sc());
        JavaPairReceiverInputDStream<String, String> messages = KafkaUtils.createStream(jsc, 
                                                                                        prop.getProperty("zookeeper.connect"),
                                                                                        prop.getProperty("group.id"), 
                                                                                        topicMap
                                                                                        );

        JavaDStream<String> lines = messages.map(new Function<Tuple2<String, String>, String>() {

                    private static final long serialVersionUID = 1L;

                    public String call(Tuple2<String, String> tuple2) {
                        return tuple2._2();
                    }
                });

我现在看到的问题是:

lines rdd 包含两个明显的行。我如何分离或找出哪些记录来自哪个主题或哪个文件。
这背后的原因是我想要为即将到来的不同主题应用不同的逻辑。但是 rdd 有时间的所有行

感谢任何建议

【问题讨论】:

    标签: apache-spark apache-kafka kafka-consumer-api


    【解决方案1】:

    您必须选择接受scala.Function1&lt;kafka.message.MessageAndMetadata&lt;K,V&gt;,R&gt; messageHandler 作为参数的createDirectStream 重载方法。然后,您只需将 messageHandler 传递给一个函数 - 获取 MessageAndMetadata 对象作为输入 - 返回实际消息和主题。

    在这里,我向您发布了用 Scala 编写的代码。您可以轻松地在 Java 中调整它:

    KafkaUtils.createDirectStream[String, String, StringDecoder, StringDecoder, (String,String)](ssc, 
            kafkaParams, 
            topicOffsetsMap, 
            (m:MessageAndMetadata[String, String])=> (m.topic,m.message()) 
            )
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2023-03-25
      • 2021-08-20
      • 2016-03-29
      • 2016-01-14
      • 2019-07-12
      • 1970-01-01
      • 2018-07-15
      相关资源
      最近更新 更多