【问题标题】:kafka KStream - topology to take n-second countskafka KStream - 采用 n 秒计数的拓扑
【发布时间】: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() 的等效操作,但我想知道为什么这段代码没有按我希望的方式工作。

【问题讨论】:

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


【解决方案1】:

我认为这与 kafka 集群设置不当有关。在将 文件描述符计数 更改为比默认值 (1024 -> 65535) 高得多的值后,这似乎符合规范。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2019-07-30
    • 2021-09-28
    • 2017-06-09
    • 1970-01-01
    • 2018-01-19
    • 2019-10-31
    • 2015-05-31
    • 2016-08-11
    相关资源
    最近更新 更多