You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink内连接出现丢数与重复数据问题排查求助

Flink窗口内连接问题:调试与排查指南

环境与业务场景

  • 运行环境:AWS Kinesis Analytics(Flink)
  • 数据源/输出:Kafka(主题12分区)
  • 核心逻辑:基于事件时间生成水位线(BoundedOutOfOrderness,延迟5秒),5秒滚动窗口,按empId字段执行内连接
  • 并行度配置:Flink并行度3

现存问题

  • 数据丢失:部分符合事件时间关联条件的记录未完成窗口内连接
  • 数据重复:部分窗口输出存在重复结果

示例数据

员工事件流

{"empName":"ted","timestamp":"0","uuid":"f2c2e48a44064d0fa8da5a3896e0e42a","empId":"23698"}
{"empName":"ted","timestamp":"1","uuid":"069f2293ad144dd38a79027068593b58","empId":"23145"}
{"empName":"john","timestamp":"2","uuid":"438c1f0b85154bf0b8e4b3ebf75947b6","empId":"23698"}
{"empName":"john","timestamp":"0","uuid":"76d1d21ed92f4a3f8e14a09e9b40a13b","empId":"23145"}
{"empName":"ted","timestamp":"0","uuid":"bbc3bad653aa44c4894d9c4d13685fba","empId":"23698"}
{"empName":"ted","timestamp":"0","uuid":"530871933d1e4443ade447adc091dcbe","empId":"23145"}
{"empName":"ted","timestamp":"1","uuid":"032d7be009cb448bb40fe5c44582cb9c","empId":"23698"}
{"empName":"john","timestamp":"1","uuid":"e5916821bd4049bab16f4dc62d4b90ea","empId":"23145"}

费用事件流

{"empId":"23698","timestamp":"0","expense":"234"}
{"empId":"23698","timestamp":"0","expense":"34"}
{"empId":"23698","timestamp":"1","expense":"234"}
{"empId":"23145","timestamp":"2","expense":"234"}
{"empId":"23698","timestamp":"2","expense":"234"}
{"empId":"23698","timestamp":"0","expense":"234"}
{"empId":"23145","timestamp":"0","expense":"234"}
{"empId":"23698","timestamp":"0","expense":"34"}
{"empId":"23145","timestamp":"1","expense":"34"}

参考代码

import java.text.SimpleDateFormat;
import java.time.Duration;
import java.time.Instant;
import java.time.LocalDateTime;
import java.time.format.DateTimeFormatter;
import java.time.temporal.ChronoUnit;
import java.util.Properties;

import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.functions.AggregateFunction;
import org.apache.flink.api.common.functions.JoinFunction;
import org.apache.flink.api.common.functions.ReduceFunction;
import org.apache.flink.api.common.state.ValueState;
import org.apache.flink.api.common.state.ValueStateDescriptor;
import org.apache.flink.api.common.typeinfo.TypeHint;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.java.functions.KeySelector;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.api.java.tuple.Tuple1;
import org.apache.flink.api.java.tuple.Tuple3;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.deser.std.StringDeserializer;
import org.apache.flink.streaming.api.TimeCharacteristic;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.datastream.KeyedStream;
import org.apache.flink.streaming.api.functions.co.CoFlatMapFunction;
import org.apache.flink.streaming.api.functions.co.KeyedCoProcessFunction;
import org.apache.flink.streaming.api.functions.co.RichCoFlatMapFunction;
import org.apache.flink.streaming.api.functions.sink.PrintSinkFunction;
import org.apache.flink.streaming.api.functions.timestamps.AscendingTimestampExtractor;
import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.assigners.TumblingProcessingTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Time;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer;
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaProducer.Semantic;
import org.apache.flink.util.Collector;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.serialization.StringSerializer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;


public class Main {

    
    private static final Logger LOG = LoggerFactory.getLogger(Main.class);
    // static String TOPIC_IN = "event_hub_all-mt-partitioned";

    static String TOPIC_ONE = "kafka_one_multi";
    static String TOPIC_TWO = "kafka_two_multi";
    static String TOPIC_OUT = "final_join_topic_multi";
    static String BOOTSTRAP_SERVER = "localhost:9092";
    
    public static void main(String[] args) {

        Producer<String> emp = new Producer<String>(BOOTSTRAP_SERVER, StringSerializer.class.getName());

        Producer<String> dept = new Producer<String>(BOOTSTRAP_SERVER, StringSerializer.class.getName());

        final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        env.setParallelism(1);

        Properties props = new Properties();
        props.put("bootstrap.servers", BOOTSTRAP_SERVER);
        props.put("client.id", "flink-example1");

        FlinkKafkaConsumer<Employee> kafkaConsumerOne = new FlinkKafkaConsumer<>(TOPIC_ONE, new EmployeeSchema(),
                props);
        
        
        LOG.info("Coming to main function");
        
        //Commenting event timestamp for watermark generation!!

        var empDebugStream = kafkaConsumerOne.assignTimestampsAndWatermarks(
                WatermarkStrategy.<Employee>forBoundedOutOfOrderness(Duration.ofSeconds(5))
                        .withTimestampAssigner((employee, timestamp) -> employee.getTimestamp().getTime())
                        .withIdleness(Duration.ofSeconds(1)));

        // for allowing Flink to handle late elements

        kafkaConsumerOne.setStartFromLatest();

        FlinkKafkaConsumer<EmployeeExpense> kafkaConsumerTwo = new FlinkKafkaConsumer<>(TOPIC_TWO,
                new DepartmentSchema(), props);
        
        
        //Commenting event timestamp for watermark generation!!

        kafkaConsumerTwo.assignTimestampsAndWatermarks(
                WatermarkStrategy.<EmployeeExpense>forBoundedOutOfOrderness(Duration.ofSeconds(5))
                        .withTimestampAssigner((employeeExpense, timestamp) -> employeeExpense.getTimestamp().getTime())
                        .withIdleness(Duration.ofSeconds(1)));

        kafkaConsumerTwo.setStartFromLatest();

//      EventSerializationSchema<EmployeeWithExpenseAggregationStats> employeeWithExpenseAggregationSerializationSchema = new EventSerializationSchema<EmployeeWithExpenseAggregationStats>(
//              TOPIC_OUT);

        EventSerializationSchema<EmployeeWithExpense> employeeWithExpenseSerializationSchema = new EventSerializationSchema<EmployeeWithExpense>(
                TOPIC_OUT);

//      FlinkKafkaProducer<EmployeeWithExpenseAggregationStats> sink = new FlinkKafkaProducer<EmployeeWithExpenseAggregationStats>(
//              TOPIC_OUT,
//              employeeWithExpenseAggregationSerializationSchema,props,
//              FlinkKafkaProducer.Semantic.AT_LEAST_ONCE);

        FlinkKafkaProducer<EmployeeWithExpense> sink = new FlinkKafkaProducer<EmployeeWithExpense>(TOPIC_OUT,
                employeeWithExpenseSerializationSchema, props, FlinkKafkaProducer.Semantic.AT_LEAST_ONCE);

        DataStream<Employee> empStream = env.addSource(kafkaConsumerOne)
                .transform("debugFilter", empDebugStream.getProducedType(), new StreamWatermarkDebugFilter<>())
                .keyBy(emps -> emps.getEmpId());

        DataStream<EmployeeExpense> expStream = env.addSource(kafkaConsumerTwo).keyBy(exps -> exps.getEmpId());

        // DataStream<EmployeeWithExpense> aggInputStream = empStream.join(expStream)

        empStream.join(expStream).where(new KeySelector<Employee, Tuple1<Integer>>() {
            /**
            * 
            */
            private static final long serialVersionUID = 1L;

            @Override
            public Tuple1<Integer> getKey(Employee value) throws Exception {
                return Tuple1.of(value.getEmpId());

            }
        }).equalTo(new KeySelector<EmployeeExpense, Tuple1<Integer>>() {

            /**
             * 
             */
            private static final long serialVersionUID = 1L;

            @Override
            public Tuple1<Integer> getKey(EmployeeExpense value) throws Exception {
                return Tuple1.of(value.getEmpId());
            }
        }).window(TumblingEventTimeWindows.of(Time.seconds(5))).allowedLateness(Time.seconds(15))
                .apply(new JoinFunction<Employee, EmployeeExpense, EmployeeWithExpense>() {

                    /**
                     * 
                     */
                    private static final long serialVersionUID = 1L;
                    
                    

                    @Override
                    public EmployeeWithExpense join(Employee first, EmployeeExpense second) throws Exception {
                        return new EmployeeWithExpense(second.getTimestamp(), first.getEmpId(), second.getExpense(),
                                first.getUuid(), LocalDateTime.now()
                                        .format(DateTimeFormatter.ofPattern("yyyy-MM-dd'T'HH:mm:ss.SSS'+0000'")));
                    }

                }).addSink(sink);

//      KeyedStream<EmployeeWithExpense, Tuple3<Integer, Integer,Long>> inputKeyedByGWNetAccountProductRTG = aggInputStream
//              .keyBy(new KeySelector<EmployeeWithExpense, Tuple3<Integer, Integer,Long>>() {
//
//                  /**
//                   * 
//                   */
//                  private static final long serialVersionUID = 1L;
//
//                  @Override
//                  public Tuple3<Integer, Integer,Long> getKey(EmployeeWithExpense value) throws Exception {
//                      return Tuple3.of(value.empId, value.expense,Instant.ofEpochMilli(value.timestamp.getTime()).truncatedTo(ChronoUnit.SECONDS).toEpochMilli());
//                  }
//              });
//
//      inputKeyedByGWNetAccountProductRTG.window(TumblingEventTimeWindows.of(Time.seconds(2)))
//              .aggregate(new EmployeeWithExpenseAggregator()).addSink(sink);

        // streamOne.print();

        // streamTwo.print();

//       DataStream<KafkaRecord> streamTwo = env.addSource(kafkaConsumerTwo);
//       
//       streamOne.connect(streamTwo).flatMap(new CoFlatMapFunction<KafkaRecord, KafkaRecord, R>() {
//      })

//
//      // Create Kafka producer from Flink API
//      Properties prodProps = new Properties();
//      prodProps.put("bootstrap.servers", BOOTSTRAP_SERVER);
//
//      FlinkKafkaProducer<KafkaRecord> kafkaProducer = 
//              
//                              new FlinkKafkaProducer<KafkaRecord>(TOPIC_OUT, 
//              
//                                                    ((record, timestamp) -> new ProducerRecord<byte[], byte[]>(TOPIC_OUT, record.key.getBytes(), record.value.getBytes())), 
//              
//                                               prodProps, 
//              
//                                     Semantic.EXACTLY_ONCE);;
//
//      DataStream<KafkaRecord> stream = env.addSource(kafkaConsumer);
//
//      stream.filter((record) -> record.value != null && !record.value.isEmpty()).keyBy(record -> record.key)
//              .timeWindow(Time.seconds(15)).allowedLateness(Time.milliseconds(500))
//              .reduce(new ReduceFunction<KafkaRecord>() {
//                  /**
//                   * 
//                   */
//                  private static final long serialVersionUID = 1L;
//                  KafkaRecord result = new KafkaRecord();
//                  @Override
//                  public KafkaRecord reduce(KafkaRecord record1, KafkaRecord record2) throws Exception 
//                  {
//                      result.key = "outKey";
//                      
//                           result.value = record1.value + " " + record2.value;
//                      
//                           return result;
//                  }
//              }).addSink(kafkaProducer);

        // produce a number as string every second
        new MessageGenerator(emp, TOPIC_ONE, "EMP").start();

        new MessageGenerator(dept, TOPIC_TWO, "EXP").start();

        // for visual topology of the pipeline. Paste the below output in
        // https://flink.apache.org/visualizer/
        // System.out.println(env.getExecutionPlan());

        // start flink
        try {
            env.execute();
            LOG.debug("Starting flink application!!");
        } catch (Exception e) {
            // TODO Auto-generated catch block
            e.printStackTrace();
        }
    }

}

技术问题与解决方案

1. 如何调试窗口触发时机?能否按窗口输出两个流的原始数据到Kafka?

调试窗口触发时机

  • 查看内置监控指标:在AWS Kinesis Analytics控制台查看Flink作业的Window Metrics,包括窗口创建、触发、关闭计数,以及水位线推进速度,直接观察窗口触发时间点。
  • 添加自定义日志:在窗口apply函数中打印窗口起止时间、当前水位线值和处理记录数,例如:
    LOG.info("窗口触发:[{}, {}],处理员工ID:{},费用记录ID:{}", 
             window.getStart(), window.getEnd(), 
             first.getEmpId(), second.getEmpId());
    

按窗口输出原始数据到Kafka

对两个流分别添加窗口处理逻辑,将每个窗口内的原始数据输出到单独Kafka主题:

// 员工流按窗口输出日志
empStream.keyBy(emps -> emps.getEmpId())
         .window(TumblingEventTimeWindows.of(Time.seconds(5)))
         .apply(new WindowFunction<Employee, String, Tuple1<Integer>, TimeWindow>() {
             @Override
             public void apply(Tuple1<Integer> key, TimeWindow window, Iterable<Employee> values, Collector<String> out) {
                 StringBuilder sb = new StringBuilder();
                 sb.append("窗口[").append(window.getStart()).append(",").append(window.getEnd()).append("],员工ID:").append(key.f0);
                 for (Employee emp : values) {
                     sb.append("\n").append(emp.toString());
                 }
                 out.collect(sb.toString());
             }
         }).addSink(new FlinkKafkaProducer<>("emp-window-logs", new SimpleStringSchema(), props));

// 费用流同理实现
expStream.keyBy(exps -> exps.getEmpId())
         .window(TumblingEventTimeWindows.of(Time.seconds(5)))
         .apply(new WindowFunction<EmployeeExpense, String, Tuple1<Integer>, TimeWindow>() {
             @Override
             public void apply(Tuple1<Integer> key, TimeWindow window, Iterable<EmployeeExpense> values, Collector<String> out) {
                 StringBuilder sb = new StringBuilder();
                 sb.append("窗口[").append(window.getStart()).append(",").append(window.getEnd()).append("],员工ID:").append(key.f0);
                 for (EmployeeExpense exp : values) {
                     sb.append("\n").append(exp.toString());
                 }
                 out.collect(sb.toString());
             }
         }).addSink(new FlinkKafkaProducer<>("exp-window-logs", new SimpleStringSchema(), props));

2. 能否将迟到记录输出到单独的Sink?

可以通过sideOutputLateData方法捕获迟到数据,输出到单独Kafka主题:

  1. 定义迟到数据标记:
private static final OutputTag<Employee> EMP_LATE_DATA = new OutputTag<Employee>("emp-late-data"){};
private static final OutputTag<EmployeeExpense> EXP_LATE_DATA = new OutputTag<EmployeeExpense>("exp-late-data"){};
  1. 在窗口配置中添加侧输出并发送到Kafka:
// 处理员工流迟到数据
DataStream<Employee> empWithLate = empStream
        .keyBy(emps -> emps.getEmpId())
        .window(TumblingEventTimeWindows.of(Time.seconds(5))
                .allowedLateness(Time.seconds(15)))
        .sideOutputLateData(EMP_LATE_DATA)
        .apply(...);

// 提取迟到数据并输出
DataStream<Employee> empLateStream = empWithLate.getSideOutput(EMP_LATE_DATA
相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.15 12:25:40