【发布时间】:2016-09-14 22:07:23
【问题描述】:
我有一个 JSON 对象流,我键入了几个值的哈希值。我希望在 n 秒(10?60?)间隔内按键计数,并使用这些值进行一些模式分析。
我的拓扑:K->aggregateByKey(n seconds)->process()
在process - init() 步骤中,我调用了ProcessorContent.schedule(60 * 1000L),希望能够调用.punctuate()。从这里我将遍历内部哈希中的值并采取相应的行动。
我看到值来自聚合步骤并点击了process() 函数,但从未调用过.punctuate()。
代码:
KStreamBuilder kStreamBuilder = new KStreamBuilder();
KStream<String, String> opxLines = kStreamBuilder.stream(TOPIC);
KStream<String, String> mapped = opxLines.map(new ReMapper());
KTable<Windowed<String>, String> ktRtDetail = mapped.aggregateByKey(
new AggregateInit(),
new OpxAggregate(),
TimeWindows.of("opx_aggregate", 60000));
ktRtDetail.toStream().process(new ProcessorSupplier<Windowed<String>, String>() {
@Override
public Processor<Windowed<String>, String> get() {
return new AggProcessor();
}
});
KafkaStreams kafkaStreams = new KafkaStreams(kStreamBuilder, streamsConfig);
kafkaStreams.start();
AggregateInit() 返回 null。
我想我可以用一个简单的计时器完成 .punctuate() 的等效操作,但我想知道为什么这段代码没有按我希望的方式工作。
【问题讨论】:
-
你的时间语义是什么(事件时间处理时间?)另见stackoverflow.com/questions/39251997/…关于标点符号的内部结构。
-
进一步探索here
标签: java apache-kafka apache-kafka-streams