【问题标题】:Cannot get a joined Kafka stream to run or output anything无法让加入的 Kafka 流运行或输出任何内容
【发布时间】:2016-08-04 15:37:21
【问题描述】:

对于下面的代码,stream1 和 stream2 都单独运行良好,我可以看到输出,但连接的流根本不记录任何内容。我感觉这与连接窗口有关,但是来自两个流的数据几乎同时进入。

val stream = builder.stream(stringSerde, byteArraySerde, "topic")

val stream1 = stream
  .filter((key, value) => somefilter(key, value))
  .through(stringSerde, byteArraySerde, "topic1")

val stream2 = stream
  .filter((key, value) => someotherfilter(key, value))
  .through(stringSerde, byteArraySerde, "topic2")

val joinedStream = stream1
  .join(stream2, (value1: Array[Byte], value2: Array[Byte]) => {
    println("wont print anything")
    return somerandomdata
  },
  JoinWindows.of("othertopic").within(10000L),
  stringSerde, byteArraySerde, byteArraySerde)

【问题讨论】:

  • 一个连接窗口是根据嵌入的记录时间戳计算的(即,除了键和值之外,每个记录中包含的元数据)。如果您打印这些时间戳以进行调试,将会有所帮助。要访问它们,您需要使用 process() -- 给定的context 对象,包含当前处理的记录的时间戳(即,为每个处理的记录更新上下文)。

标签: java scala apache-kafka apache-kafka-streams


【解决方案1】:

两个主题的键不应该相同才能加入吗?

我认为 Javadoc 解释了这一点: https://kafka.apache.org/0102/javadoc/org/apache/kafka/streams/kstream/JoinWindows.html

这也可能是一个有趣的阅读: https://cwiki.apache.org/confluence/display/KAFKA/Kafka+Streams+Join+Semantics

【讨论】:

    猜你喜欢
    • 2010-12-15
    • 1970-01-01
    • 1970-01-01
    • 2015-06-30
    • 2012-12-31
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多