【问题标题】:KTable Reduce function does not honor windowingKTable Reduce 函数不支持窗口化
【发布时间】:2018-02-21 12:26:34
【问题描述】:

要求:- 我们需要合并所有具有相同orderid的消息,并对合并后的消息进行后续操作。

说明:- 下面的代码 sn-p 尝试捕获从特定租户收到的所有订单消息,并在等待特定时间段后尝试合并为单个订单消息 它做了以下事情

  1. 基于 OrderId 的重新分区消息。因此,每条订单消息都将以tenantId 和 groupId 作为其键
  2. 执行 groupby 键操作,然后执行窗口操作 2 分钟
  3. 一旦窗口完成,就会执行缩减操作。
  4. 再次将 Ktable 转换为流式返回,然后将其输出发送到另一个 kafka 主题

预期输出:- 如果在窗口期内发送了 5 条具有相同订单 ID 的消息。预计最终的kafka topic应该只有一条消息,并且是最后一个reduce操作消息。

实际输出:- 看到所有 5 条消息,表明在调用 reduce 操作之前没有发生窗口化。在 kafka 中看到的所有消息都在收到每条消息时都进行了适当的 reduce 操作。

查询:- 在 kafka 流库版本 0.11.0.0 中,reduce 函数用于接受 timewindow 作为其参数。我看到这在 kafka 流版本 1.0.0 中已被弃用。在下面的代码中完成的窗口化,是否正确?较新版本的 kafka 流库 1.0.0 是否支持窗口化?如果是这样,那么下面的sn-p代码有什么可以改进的吗?

        String orderMsgTopic = "sampleordertopic";

        JsonSerializer<OrderMsg> orderMsgJSONSerialiser = new JsonSerializer<>();
        JsonDeserializer<OrderMsg> orderMsgJSONDeSerialiser = new JsonDeserializer<>(OrderMsg.class);

        Serde<OrderMsg> orderMsgSerde = Serdes.serdeFrom(orderMsgJSONSerialiser,orderMsgJSONDeSerialiser);



        KStream<String, OrderMsg> orderMsgStream = this.builder.stream(orderMsgTopic, Consumed.with(Serdes.ByteArray(), orderMsgSerde))
                                                                .map(new KeyValueMapper<byte[], OrderMsg, KeyValue<? extends String, ? extends OrderMsg>>() {
                                                                    @Override
                                                                    public KeyValue<? extends String, ? extends OrderMsg> apply(byte[] byteArr, OrderMsg value) {
                                                                        TenantIdMessageTypeDeserializer deserializer = new TenantIdMessageTypeDeserializer();
                                                                        TenantIdMessageType tenantIdMessageType = deserializer.deserialize(orderMsgTopic, byteArr);
                                                                        String newTenantOrderKey = null;
                                                                        if ((tenantIdMessageType != null) && (tenantIdMessageType.getMessageType() == 1)) {
                                                                            Long tenantId = tenantIdMessageType.getTenantId();
                                                                            newTenantOrderKey = tenantId.toString() + value.getOrderKey();
                                                                        } else {
                                                                            newTenantOrderKey = value.getOrderKey();
                                                                        }
                                                                        return new KeyValue<String, OrderMsg>(newTenantOrderKey, value);
                                                                    }
                                                                });



        final KTable<Windowed<String>, OrderMsg> orderGrouping = orderMsgStream.groupByKey(Serialized.with(Serdes.String(), orderMsgSerde))
                                                                                .windowedBy(TimeWindows.of(windowTime).advanceBy(windowTime))
                                                                                .reduce(new OrderMsgReducer());


        orderGrouping.toStream().map(new KeyValueMapper<Windowed<String>, OrderMsg, KeyValue<String, OrderMsg>>() {
                                                                    @Override
                                                                    public KeyValue<String, OrderMsg> apply(Windowed<String> key, OrderMsg value) {
                                                                        return new KeyValue<String, OrderMsg>(key.key(), value);
                                                                    }
                                                                }).to("newone11", Produced.with(Serdes.String(), orderMsgSerde));

【问题讨论】:

标签: apache-kafka apache-kafka-streams


【解决方案1】:

我意识到我已将 StreamsConfig.CACHE_MAX_BYTES_BUFFERING_CONFIG 设置为 0,并将默认提交间隔设置为 1000 毫秒。更改此值在某种程度上可以帮助我使窗口正常工作

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-06-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-02-10
    • 1970-01-01
    • 2012-10-21
    相关资源
    最近更新 更多