【问题标题】:Joining streams Flink doesn't work with Kafka consumer加入流 Flink 不适用于 Kafka 消费者
【发布时间】:2022-06-10 22:18:25
【问题描述】:

我正在尝试加入两个流,一个来自数据收集,一个来自 Kafka。

代码 sn-p

public static void main(String[] args) {
        KafkaSource<JsonNode> kafkaSource = ...
        
        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        
        // Kafka messages : {"name": "John"} 
        final DataStream<JsonNode> dataStream1 = env.fromSource(kafkaSource, waterMark(), "Kafka").rebalance()
                .assignTimestampsAndWatermarks(waterMark());
        
        final DataStream<String> dataStream2 = env.fromElements("John", "Zbe", "Abe")
                .assignTimestampsAndWatermarks(waterMark());
        
        dataStream1
            .join(dataStream2)
            .where(new KeySelector<JsonNode, String>() {
    
                @Override
                public String getKey(JsonNode value) throws Exception {
                    return value.get("name").asText();
                }
            })
            .equalTo(new KeySelector<String, String>() {

                @Override
                public String getKey(String value) throws Exception {
                    return value;
                }
            })
            .window(SlidingEventTimeWindows.of(Time.minutes(50) /* size */, Time.minutes(10) /* slide */))
            .apply(new JoinFunction<JsonNode, String, String>() {
    
                @Override
                public String join(JsonNode first, String second) throws Exception {

                    return first+" "+second;
                }
            }).print();
            
            env.execute();
    }

水印

private static <T>  WatermarkStrategy<T> waterMark() {
        return new WatermarkStrategy<T>() {

            @Override
            public WatermarkGenerator<T> createWatermarkGenerator(
                    org.apache.flink.api.common.eventtime.WatermarkGeneratorSupplier.Context context) {
                return new AscendingTimestampsWatermarks<>();
            }
            
            @Override
            public TimestampAssigner<T> createTimestampAssigner(TimestampAssignerSupplier.Context context) {
                return (event, timestamp) -> System.currentTimeMillis();
            }
            
        };
    }

运行 sn-p 代码后,输出中没有任何合并数据。我是不是哪里出错了?

Apache flink 版本:1.13.2

【问题讨论】:

    标签: apache-kafka apache-flink


    【解决方案1】:

    问题可能与水印有关。由于您没有使用基于事件时间的时间戳,请尝试将 SlidingEventTimeWindows 更改为 SlidingProcessingTimeWindows 并查看它是否会产生结果。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2017-08-18
      • 2018-10-15
      • 1970-01-01
      • 2016-12-03
      • 1970-01-01
      • 1970-01-01
      • 2018-12-31
      • 2018-12-03
      相关资源
      最近更新 更多