【发布时间】: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<Windowed<String>, Long> 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