通过FlinkCDCKafka对变更数据排序,核心思路是按键分组,利用窗口函数对每组数据排序。具体示例中,程序从Kafka中读取数据,解析键值,按照键分组,在五分钟的滚动窗口内按值进行排序并输出,实现有序处理,保证对于同一主键的变更数据按时间顺序输出,减少数据乱序问题,提升下游应用准确性,适用于实时。
在实际的流处理场景中,经常需要对从 Kafka 捕获的变更数据进行排序。那么,具体该怎么做呢?核心思路其实很直接:根据变更数据的键进行分组,然后利用 Flink 的窗口函数对每个分组内的数据依次排序。下面通过一个完整的示例来拆解每一步。

长期稳定更新的攒劲资源: >>>点此立即查看<<<
先把 Flink CDC Kafka 依赖加到项目里。如果使用 Maven,在 pom.xml 中添加:
<dependency><groupId>com.ververicagroupId><artifactId>flink-connector-kafka-cdc_2.11artifactId><version>1.14.0version>dependency>创建一个 Flink 程序,用 KafkaSourceBuilder 从 Kafka 读取变更数据。假设主题名为 my-topic,且 Kafka 连接器已配置完成。
import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;import org.apache.flink.streaming.util.serialization.SimpleStringSchema;FlinkKafkaConsumer kafkaConsumer = new FlinkKafkaConsumer<>("my-topic", new SimpleStringSchema(), properties);DataStream stream = env.addSource(kafkaConsumer); 解析变更数据,提取出键和值。这里假设数据格式为 JSON,且键和值用逗号分隔。完整的主类代码如下:
import org.apache.flink.api.common.typeinfo.TypeInformation;import org.apache.flink.api.java.typeutils.TypeExtractor;import org.apache.flink.streaming.api.datastream.DataStream;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;import org.apache.flink.streaming.util.serialization.SimpleStringSchema;import java.util.Properties;public class FlinkCdcKafkaSort {public static void main(String[] args) throws Exception {final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();Properties properties = new Properties();properties.setProperty("bootstrap.servers", "localhost:9092");properties.setProperty("group.id", "flink-cdc-kafka-sort");FlinkKafkaConsumer kafkaConsumer = new FlinkKafkaConsumer<>("my-topic", new SimpleStringSchema(), properties);DataStream stream = env.addSource(kafkaConsumer);DataStream changeRecords = stream.map(new ChangeRecordParser()).keyBy(ChangeRecord::getKey).window(TumblingEventTimeWindows.of(Time.minutes(5))).apply(new SortFunction());changeRecords.print();env.execute("Flink CDC Kafka Sort");}} 创建一个 ChangeRecordParser 类,用于解析变更数据,实现 MapFunction 接口。
import org.apache.flink.api.common.typeinfo.TypeInformation;import org.apache.flink.api.java.typeutils.TypeExtractor;import org.apache.flink.streaming.api.datastream.DataStream;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;import org.apache.flink.streaming.util.serialization.SimpleStringSchema;import java.util.Properties;public class FlinkCdcKafkaSort {public static void main(String[] args) throws Exception {// ... 省略其他代码 ...}public static class ChangeRecordParser implements MapFunction {@Overridepublic ChangeRecord map(String value) throws Exception {// 解析变更数据,提取键和值String[] parts = value.split(",");String key = parts[0];String value = parts[1];// 创建并返回 ChangeRecord 对象return new ChangeRecord(key, value);}}} 创建一个 ChangeRecord 类,表示变更记录,实现 Serializable 接口,包含键和值属性。
import java.io.Serializable;public class ChangeRecord implements Serializable {private String key;private String value;public ChangeRecord(String key, String value) {this.key = key;this.value = value;}public String getKey() {return key;}public String getValue() {return value;}}创建一个 SortFunction 类,实现 WindowFunction 接口,对窗口内的记录按值排序。
import org.apache.flink.api.common.typeinfo.TypeInformation;import org.apache.flink.api.java.typeutils.TypeExtractor;import org.apache.flink.streaming.api.datastream.DataStream;import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;import org.apache.flink.streaming.connectors.kafka.FlinkKafkaConsumer;import org.apache.flink.streaming.util.serialization.SimpleStringSchema;import org.apache.flink.streaming.api.windowing.time.Time;import org.apache.flink.streaming.api.windowing.windows.TimeWindow;import org.apache.flink.streaming.api.functions.windowing.WindowFunction;import org.apache.flink.util.Collector;import java.util.List;public class FlinkCdcKafkaSort {// ... 省略其他代码 ...public static class SortFunction extends WindowFunction {@Overridepublic void apply(String key, TimeWindow window, Iterable input, Collector out) {List sortedRecords = input.stream().sorted((record1, record2) -> record1.getValue().compareTo(record2.getValue())).collect(Collectors.toList());for (ChangeRecord record : sortedRecords) {out.collect(new SortedChangeRecord(record.getKey(), record.getValue()));}}}} 最后,创建一个 SortedChangeRecord 类,表示已排序的变更记录,同样实现 Serializable,包含键和值。
import java.io.Serializable;public class SortedChangeRecord implements Serializable {private String key;private String value;public SortedChangeRecord(String key, String value) {this.key = key;this.value = value;}public String getKey() {return key;}public String getValue() {return value;}}运行这个 Flink 程序后,它会持续从 Kafka 读取变更数据,按 key 分组,在每个 5 分钟的滚动窗口内对每个组的数据按值排序,最终输出排好序的记录。整个过程清晰、可复用,适合大多数 CDC 排序需求。
侠游戏发布此文仅为了传递信息,不代表侠游戏网站认同其观点或证实其描述