【问题标题】:How to capture kafka records that don't match the condition of the kafka stream join?如何捕获不符合kafka流连接条件的kafka记录?
【发布时间】:2019-11-03 03:27:05
【问题描述】:

我正在通过将 kstream 与 ktable 连接起来来丰富数据。 kstream 包含车辆发送的消息,ktable 包含车辆数据。 我遇到的问题是我想从流中捕获表中没有相应连接键的消息。 Kafka 流静默地跳过它们没有连接匹配的记录。 有没有办法将这些记录发送到不同的主题,以便以后处理?

StreamsBuilder builder = new StreamsBuilder();
        final KTable<String, VinMappingInfo> vinMappingTable = builder.table(vinInfoTopic, Consumed.with(Serdes.String(), valueSerde));
        KStream<String, VehicleMessage> vehicleStream = builder.stream(sourceTopic);
        vehicleStream.join(vinMappingTable, (vehicleMsg, vinInfo) -> {
            log.info("joining {} with vin info {}", vehicleMsg.getPayload().getId(), vinInfo.data.vin);
            vehicleMsg.setVin(vinInfo.data.vin);
            return vehicleMsg;
        }, Joined.with(null, null, valueSerde))
                .to(destinationTopic);

        final Topology topology = builder.build();
        log.info("The topology of connected processor nodes: \n {}", topology.describe());
        KafkaStreams streams = new KafkaStreams(topology, config);
        streams.cleanUp();
        streams.start();

【问题讨论】:

    标签: join apache-kafka apache-kafka-streams


    【解决方案1】:

    您可以使用左连接:

    stream.leftJoin(table,...);
    

    这确保来自输入流的所有记录都在输出流中。在这种情况下,ValueJoiner 将与 apply(streamValue, null) 一起调用。

    【讨论】:

    • 嗨,Matthias J. Sax,我正在尝试获取在窗口期到期后 KStream 到 KStream 加入期间未处理的事件。这是为了在加入流数据稍后到达时重试。
    • 这不容易。
    • @matthias-j-sax 有什么办法可以识别???因为,我现在使用 KStream join with window,一个流可能会延迟 1 或 2 或 3 天。我不想有这么大的窗口/宽限期。因为,我的应用程序每天将有大约 5000 万个事件的高数据流。所以我想识别非加入消息并以不同的方式处理。感谢您的帮助,因此请解决此问题。
    • 您可以使用上游 transformValue 手动跟踪“流时间”和“分支”记录,这些记录将在加入前删除。 “流时间”只是看到的最大事件时间戳,因此很容易计算——但是,您可能需要一个状态存储来使跟踪容错(直到 KIP-622 登陆:cwiki.apache.org/confluence/display/KAFKA/…
    • Matthias,你是对的 - 事实证明我没有仔细阅读这个问题。我对在流/流加入场景中报告未加入的消息很感兴趣。现在我了解您的建议,因为它与流/表连接有关。我将听从您的建议并开始新的讨论(可能在论坛中)!非常感谢!
    猜你喜欢
    • 1970-01-01
    • 2021-08-23
    • 2019-03-10
    • 1970-01-01
    • 2020-08-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多