【问题标题】:Buffer messages in stream data for a given messageId在给定 messageId 的流数据中缓冲消息
【发布时间】:2018-02-13 02:59:44
【问题描述】:

用例:我有具有 messageId 的消息,多条消息可以具有相同的消息 ID,这些消息存在于由 messageId 分区的流式管道(如 kafka)中,因此我确保所有具有相同 messageId 的消息都将进入相同的分区。

所以我需要编写一个作业,它应该将消息缓冲一段时间(比如说 1 分钟),然后将具有相同 messageId 的所有消息合并为单个大消息。

我认为可以使用 spark Datasets 和 spark sql(或其他东西?)来完成。但是我找不到任何关于如何为给定消息 ID 存储消息一段时间然后对这些消息进行聚合的示例/文档。

【问题讨论】:

  • 你在想什么样的聚合?你想要一个聚合值(比如一个总和),还是想要一个消息的消息?
  • 消息的消息,假设您有 10 条消息具有相同的消息 ID,我的结果应该是 1 条大消息,其中包含所有 10 条消息。希望清除。

标签: apache-kafka streaming spark-streaming buffering apache-samza


【解决方案1】:

我认为您正在寻找的是Spark Streaming。 Spark 有一个 Kafka Connector 可以链接到 Spark Streaming Context。

这是一个非常基本的示例,它将在 1 分钟的时间间隔内为给定主题集中的所有消息创建一个 RDD,然后按消息 id 字段对它们进行分组(您的值序列化程序必须公开这样的 getMessageId 方法,当然)。

SparkConf conf = new SparkConf().setAppName(appName);
JavaStreamingContext ssc = new JavaStreamingContext(conf, Durations.minutes(1));

Map<String, Object> params = new HashMap<String, Object>() {{
    put("bootstrap.servers", kafkaServers);
    put("key.deserializer", kafkaKeyDeserializer);
    put("value.deserializer", kafkaValueDeserializer);
}};

List<String> topics = new ArrayList<String>() {{
    // Add Topics
}};

JavaInputDStream<ConsumerRecord<String, String>> stream =
    KafkaUtils.createDirectStream(ssc,
        LocationStrategies.PreferConsistent(),
        ConsumerStrategies.<String, String>Subscribe(topics, params)
    );

stream.foreachRDD(rdd -> rdd.groupBy(record -> record.value().getMessageId()));

ssc.start();
ssc.awaitTermination(); 

还有其他几种方法可以在流式 API 中对消息进行分组。查看文档以获取更多示例。

【讨论】:

猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2011-08-05
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2017-10-19
  • 1970-01-01
相关资源
最近更新 更多