Flink CEP+CDC 实战:实时数据捕获与复杂事件处理
·
一、前言
在实时计算领域,Flink CDC(Change Data Capture)负责实时捕获数据库全量 + 增量数据变更,Flink CEP(Complex Event Processing)专注复杂事件模式匹配,二者结合可快速构建 “数据实时同步 + 异常行为检测” 的端到端实时链路,广泛应用于金融风控、电商实时营销、物联网监控等场景。
二、核心原理与应用场景
2.1 Flink CDC 核心原理
CDC 是变更数据捕获技术,Flink CDC 内置Debezium 引擎,直接读取数据库日志(如 MySQL binlog、PostgreSQL WAL),无需中间件,实现全量快照 + 增量实时同步,核心特点:
- 低侵入:不修改业务表,无数据库压力;
- 高实时:毫秒级捕获 INSERT/UPDATE/DELETE;
- Exactly-Once:基于 Flink Checkpoint 保证数据一致性;
- 多源支持:MySQL、PostgreSQL、Oracle、MongoDB 等。

2.2 Flink CEP 核心原理
CEP 是复杂事件处理,用于在流数据中匹配预设事件模式(如 “10 分钟内 3 次登录失败”),核心 API:
- Pattern:定义事件序列模式(begin/next/followedBy/where/within);
- PatternStream:应用模式到数据流;
- Select/FlatSelect:提取匹配结果。

2.3 组合应用场景
- 金融风控:CDC 捕获支付流水,CEP 检测 “10 分钟内异地多笔支付” 盗刷行为;
- 电商实时营销:CDC 捕获用户浏览 / 下单数据,CEP 匹配 “5 分钟内浏览同商品 3 次”,触发优惠券推送;
- 运维监控:CDC 捕获服务器日志库数据,CEP 匹配 “连续 3 次接口超时”,触发告警。
三、环境准备
3.1 核心依赖(Maven)
<properties>
<flink.version>1.17.0</flink.version>
<java.version>1.8</java.version>
<scala.binary.version>2.12</scala.binary.version>
</properties>
<dependencies>
<!-- Flink 核心 -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-java</artifactId>
<version>${flink.version}</version>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java</artifactId>
<version>${flink.version}</version>
</dependency>
<!-- Flink CDC MySQL -->
<dependency>
<groupId>com.ververica</groupId>
<artifactId>flink-connector-mysql-cdc</artifactId>
<version>2.4.0</version>
</dependency>
<!-- Flink CEP -->
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-cep</artifactId>
<version>${flink.version}</version>
</dependency>
<!-- 日志 -->
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-simple</artifactId>
<version>1.7.36</version>
<scope>test</scope>
</dependency>
</dependencies>
3.2 MySQL 配置(开启 binlog)
修改my.cnf(Linux)/my.ini(Windows),重启 MySQL:
[mysqld]
server-id=1
log-bin=mysql-bin
binlog-format=ROW
binlog-do-db=flink_cdc_db # 同步的数据库名
3.3 测试表结构
CREATE DATABASE IF NOT EXISTS flink_cdc_db;
USE flink_cdc_db;
CREATE TABLE pay_event (
id BIGINT PRIMARY KEY AUTO_INCREMENT,
user_id INT,
amount DECIMAL(10,2),
location VARCHAR(50),
ts VARCHAR(20) COMMENT 'yyyy-MM-dd HH:mm:ss'
);
四、完整实战:CDC 采集 + CEP 支付风控
4.1 业务需求
实时捕获 MySQL 支付表pay_event的数据变更,通过 CEP 匹配“同一用户 10 分钟内,在 2 个不同地点发生≥2 笔支付”的盗刷风险,输出告警信息。
4.2 步骤 1:定义 POJO 类(PayEvent)
import java.util.Date;
public class PayEvent {
private Long id;
private Integer userId;
private Double amount;
private String location;
private String ts; // 时间字符串 yyyy-MM-dd HH:mm:ss
// 无参构造(Flink序列化必需)
public PayEvent() {}
// getter/setter
public Long getId() { return id; }
public void setId(Long id) { this.id = id; }
public Integer getUserId() { return userId; }
public void setUserId(Integer userId) { this.userId = userId; }
public Double getAmount() { return amount; }
public void setAmount(Double amount) { this.amount = amount; }
public String getLocation() { return location; }
public void setLocation(String location) { this.location = location; }
public String getTs() { return ts; }
public void setTs(String ts) { this.ts = ts; }
@Override
public String toString() {
return "PayEvent{userId=" + userId + ", amount=" + amount + ", location='" + location + "', ts='" + ts + "'}";
}
}
4.3 步骤 2:Flink CDC 采集 MySQL 数据
import com.ververica.cdc.connectors.mysql.source.MySqlSource;
import com.ververica.cdc.debezium.JsonDebeziumDeserializationSchema;
import org.apache.flink.api.common.eventtime.SerializableTimestampAssigner;
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.ProcessFunction;
import org.apache.flink.util.Collector;
import org.json.JSONObject;
import java.text.ParseException;
import java.text.SimpleDateFormat;
import java.time.Duration;
import java.util.Date;
public class CdcSourceDemo {
public static void main(String[] args) throws Exception {
// 1. 初始化流环境
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(1); // 本地测试设为1
env.enableCheckpoint(3000); // 开启Checkpoint,3s一次
// 2. 构建MySQL CDC Source
MySqlSource<String> mySqlSource = MySqlSource.<String>builder()
.hostname("localhost")
.port(3306)
.databaseList("flink_cdc_db") // 库名
.tableList("flink_cdc_db.pay_event") // 表名
.username("root")
.password("123456")
.deserializer(new JsonDebeziumDeserializationSchema()) // 序列化为JSON
.build();
// 3. 读取CDC数据流
DataStream<String> cdcStream = env.fromSource(mySqlSource, WatermarkStrategy.noWatermarks(), "MySQL-CDC-Source");
// 4. 解析JSON为PayEvent对象(过滤删除操作)
DataStream<PayEvent> payEventStream = cdcStream.process(new ProcessFunction<String, PayEvent>() {
private static final long serialVersionUID = 1L;
private final SimpleDateFormat sdf = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
@Override
public void processElement(String value, Context ctx, Collector<PayEvent> out) throws Exception {
try {
JSONObject json = new JSONObject(value);
String op = json.getString("op"); // 操作类型:c=新增,u=更新,d=删除
if ("d".equals(op)) return; // 过滤删除
JSONObject after = json.getJSONObject("after"); // 变更后数据
PayEvent event = new PayEvent();
event.setId(after.getLong("id"));
event.setUserId(after.getInt("user_id"));
event.setAmount(after.getDouble("amount"));
event.setLocation(after.getString("location"));
event.setTs(after.getString("ts"));
out.collect(event);
} catch (Exception e) {
e.printStackTrace();
}
}
});
// 5. 分配水位线(事件时间,处理乱序)
DataStream<PayEvent> watermarkStream = payEventStream.assignTimestampsAndWatermarks(
WatermarkStrategy.<PayEvent>forBoundedOutOfOrderness(Duration.ofSeconds(3)) // 允许3s乱序
.withTimestampAssigner(new SerializableTimestampAssigner<PayEvent>() {
@Override
public long extractTimestamp(PayEvent element, long recordTimestamp) {
try {
Date date = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").parse(element.getTs());
return date.getTime();
} catch (ParseException e) {
e.printStackTrace();
return 0;
}
}
})
);
// 打印原始数据(测试用)
watermarkStream.print("CDC原始数据:");
// 6. 调用CEP逻辑
CepRiskDetection.detectRisk(watermarkStream);
// 7. 执行任务
env.execute("Flink-CDC-CEP-Pay-Risk-Demo");
}
}
4.4 步骤 3:Flink CEP 风控逻辑实现
import org.apache.flink.cep.CEP;
import org.apache.flink.cep.PatternSelectFunction;
import org.apache.flink.cep.PatternStream;
import org.apache.flink.cep.pattern.Pattern;
import org.apache.flink.cep.pattern.conditions.SimpleCondition;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.windowing.time.Time;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
public class CepRiskDetection {
public static void detectRisk(DataStream<PayEvent> payStream) {
// 1. 定义CEP模式:同一用户10分钟内,≥2笔支付,且地点不同
Pattern<PayEvent, ?> riskPattern = Pattern.<PayEvent>begin("firstPay")
.where(new SimpleCondition<PayEvent>() {
@Override
public boolean filter(PayEvent value) throws Exception {
return value.getAmount() > 0; // 过滤有效支付
}
})
.followedBy("secondPay") // 宽松匹配(中间可夹杂其他事件)
.where(new SimpleCondition<PayEvent>() {
@Override
public boolean filter(PayEvent value) throws Exception {
return value.getAmount() > 0;
}
})
.within(Time.minutes(10)); // 10分钟内
// 2. 按用户分组,应用模式
PatternStream<PayEvent> patternStream = CEP.pattern(
payStream.keyBy(PayEvent::getUserId), // 按用户ID分组
riskPattern
);
// 3. 提取匹配结果,生成告警
DataStream<String> alertStream = patternStream.select(new PatternSelectFunction<PayEvent, String>() {
@Override
public String select(Map<String, List<PayEvent>> map) throws Exception {
PayEvent first = map.get("firstPay").get(0);
PayEvent second = map.get("secondPay").get(0);
// 校验地点不同
if (!first.getLocation().equals(second.getLocation())) {
return "【风控告警】用户" + first.getUserId() + "疑似盗刷!"
+ "10分钟内异地支付:" + first.getLocation() + "→" + second.getLocation()
+ ",金额:" + first.getAmount() + "→" + second.getAmount();
}
return null;
}
});
// 4. 打印告警结果
alertStream.print("CEP风控告警:");
}
}
4.5 测试验证
- 启动
CdcSourceDemo; - 向 MySQL
pay_event表插入测试数据:
-- 第1笔:用户1001,北京,2026-06-11 10:00:00
INSERT INTO pay_event (user_id, amount, location, ts) VALUES (1001, 500, '北京', '2026-06-11 10:00:00');
-- 第2笔:用户1001,上海,2026-06-11 10:05:00(触发告警)
INSERT INTO pay_event (user_id, amount, location, ts) VALUES (1001, 800, '上海', '2026-06-11 10:05:00');
控制台输出:
CDC原始数据:PayEvent{userId=1001, amount=500.0, location='北京', ts='2026-06-11 10:00:00'}
CDC原始数据:PayEvent{userId=1001, amount=800.0, location='上海', ts='2026-06-11 10:05:00'}
CEP风控告警:【风控告警】用户1001疑似盗刷!10分钟内异地支付:北京→上海,金额:500.0→800.0
五、高频踩坑总结
5.1 CDC 常见问题
- 无法连接 MySQL:检查
hostname/port/账号密码,确认 MySQL 远程访问权限; - 无数据输出:
- 检查 binlog 是否开启,
server-id是否唯一; - 确认
databaseList/tableList配置正确; - 更换
group.id(避免偏移量卡住)。
- 检查 binlog 是否开启,
- JSON 解析异常:确保
deserializer为JsonDebeziumDeserializationSchema,处理op为d的删除数据。
5.2 CEP 常见问题
- 无匹配结果:
- 水位线问题:时间解析错误、乱序时间过长(调大
Duration); - 模式过严:
next改为followedBy,放宽匹配规则; within时间不足:确保事件在指定时间内。
- 水位线问题:时间解析错误、乱序时间过长(调大
- 水位线卡住:检查时间戳解析是否正确,避免返回 0 或异常值。
- 并行度问题:本地测试设
env.setParallelism(1),避免数据分区导致无输出。
5.3 组合优化建议
- CDC+Kafka 解耦:CDC 数据写入 Kafka,CEP 消费 Kafka,提升稳定性;
- 状态后端配置:生产环境启用 RocksDB 状态后端,支持大状态;
- CEP 模式优化:复杂模式拆分,减少状态存储压力。
六、总结
Flink CDC 解决了数据实时采集的痛点,Flink CEP 提供了复杂事件匹配的能力,二者结合可快速搭建低延迟、高可靠的实时风控 / 营销系统。
更多推荐



所有评论(0)