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主题:
- 定义迟到数据标记:
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"){};
- 在窗口配置中添加侧输出并发送到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
相关产品推荐
相关产品推荐

