一、前言

在实时计算领域,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 测试验证

  1. 启动CdcSourceDemo
  2. 向 MySQLpay_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 常见问题

  1. 无法连接 MySQL:检查hostname/port/账号密码,确认 MySQL 远程访问权限;
  2. 无数据输出
    • 检查 binlog 是否开启,server-id是否唯一;
    • 确认databaseList/tableList配置正确;
    • 更换group.id(避免偏移量卡住)。
  3. JSON 解析异常:确保deserializerJsonDebeziumDeserializationSchema,处理opd的删除数据。

5.2 CEP 常见问题

  1. 无匹配结果
    • 水位线问题:时间解析错误、乱序时间过长(调大Duration);
    • 模式过严:next改为followedBy,放宽匹配规则;
    • within时间不足:确保事件在指定时间内。
  2. 水位线卡住:检查时间戳解析是否正确,避免返回 0 或异常值。
  3. 并行度问题:本地测试设env.setParallelism(1),避免数据分区导致无输出。

5.3 组合优化建议

  • CDC+Kafka 解耦:CDC 数据写入 Kafka,CEP 消费 Kafka,提升稳定性;
  • 状态后端配置:生产环境启用 RocksDB 状态后端,支持大状态;
  • CEP 模式优化:复杂模式拆分,减少状态存储压力。

六、总结

Flink CDC 解决了数据实时采集的痛点,Flink CEP 提供了复杂事件匹配的能力,二者结合可快速搭建低延迟、高可靠的实时风控 / 营销系统。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐