【问题标题】:kafka stream windowed count output unreadablekafka流窗口计数输出不可读
【发布时间】:2018-05-01 07:02:08
【问题描述】:

我正在尝试使用字数统计示例进行窗口计数。它工作正常,只是输出部分不可读。

代码:

    StringSerializer stringSerializer = new StringSerializer();
    StringDeserializer stringDeserializer = new StringDeserializer();
    WindowedSerializer<String> windowedSerializer = new WindowedSerializer<>(stringSerializer);
    WindowedDeserializer<String> windowedDeserializer = new WindowedDeserializer<>(stringDeserializer);
    Serde<Windowed<String>> windowedSerde = Serdes.serdeFrom(windowedSerializer, windowedDeserializer);

    TimeWindows window = TimeWindows.of(TimeUnit.MINUTES.toMillis(1)).advanceBy(TimeUnit.MINUTES.toMillis(1));

    KStream<String, String> textLines = builder.stream("streams-plaintext-input");
    KTable<Windowed<String>, Long> wordCounts = textLines
        .flatMapValues(textLine -> Arrays.asList(textLine.toLowerCase().split("\\W+")))
        .groupBy((key, word) -> word)
        .windowedBy(window)
        .count(Materialized.<String, Long, WindowStore<Bytes, byte[]>>as("counts-store"));
    wordCounts.toStream().to("streams-plaintext-output", Produced.with(windowedSerde, Serdes.Long()));

    KafkaStreams streams = new KafkaStreams(builder.build(), config);
    streams.start();

输出:

kafka c[??   1
yaya c[??    1
kafka c[??   2

我猜不可读的部分可能是 Windows 持续时间。 我该怎么做才能让它可读?

编辑:

尝试使用windowedSerde打印输出:

    KStream<Windowed<String>, Long> output = builder.stream("streams-plaintext-output");
    output.print(windowedSerde, Serdes.Long());

还是不行。

【问题讨论】:

  • 原代码中没有print()声明?您是否尝试使用 bin/kafka-console-consumer.sh 阅读?对于 EDIT 部分,您需要在 builder.stream() 运算符中指定 windowedSerde
  • 是的,之前我使用kafka-conole-consumer阅读主题。如果我修改为KStream&lt;Windowed&lt;String&gt;, Long&gt; output = builder.stream("streams-plaintext-output", Consumed.with( windowedSerde, Serdes.Long()));,它将显示Deserialization exception handler is set to fail upon a deserialization error
  • 您需要在builder.stream(...)中指定windowedSerde

标签: apache-kafka apache-kafka-streams


【解决方案1】:

从主题中读取时,您需要使用适合用于生成主题的序列化器的反序列化器。在这种情况下,您需要使用已经在构建的 windowDeserializer,如下所示:

WindowedDeserializer<String> windowedDeserializer = new WindowedDeserializer<>(stringDeserializer);

【讨论】:

  • 感谢您的帮助。我尝试使用 windowedSerde 来阅读主题,但仍然没有运气。附在这篇文章中的代码。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2016-12-04
  • 1970-01-01
  • 2020-09-13
  • 1970-01-01
  • 2023-02-15
相关资源
最近更新 更多