摘要
postgres是通过复制槽来管理日志的,从而确保主从的数据一致,因此flinkcdc用到postgres的复制槽,授权方面和oracle区别蛮大的。

1. Flink CDC同步PG到Doris的的作业全过程

PG主库:业务操作→PG写盘(内存WAL缓冲区)→PG刷盘(硬盘WAL文件)→复制槽标记可拉取
→
Flink CDC:拉取WAL→CDC写盘(TaskManager内存缓冲区)→CDC重放(解析WAL为增删改事件)→CDC刷盘(状态后端存位点+事件)→更新
复制槽位点→
Doris:接收事件→落地数据→Checkpoint完成

2.postgres授权flinkcdc读取日志

1)创建具有流复制权限的用户flinkcdc

create user flinkcdc login replication encrypted password '123456';


2) 给用户flinkcdc赋予数据库连接权限
grant connect on database postgres to flinkcdc;

3)把当前库所有表查询权限赋给用户flinkcdc
grant select on all tables in schema public to flinkcdc;

4) 发布表(要抽取的表都要进行发布)
alter publication dbz_publication add table schema_name.表名;

5) 更改复制标识包含更新和删除之前值(要抽取的表,如果有分区,只需要对分区进行操作,主表可以不操作,分区要有继承功能,这样后面不用每次建分区去授权)
alter table 表名 replica identity full;

4)和5)可以确保flinkcdc捕捉到删除、修改的数据

6) 要注意的是postgres如果主表或分表有分区,flinkcdc只能捕捉到分区的数据,捕捉不到主表或分表的数据,如果分区的名称有规律,可以用正则匹配,匹配到的数据,都写入到目标表,目标表可以是其他数据库的表名(理论上和抽取的主表名称一样,或者加上ods这类的前缀)

3. 排查CDC问题

 需要用到的查询命令(这些都是常用的命令)
1) 查询哪些表已经发布
select * from pg_publication_tables;

 2)查看复制标识(为f标识说明设置成功,能抽取到表变更的日志)
select relreplident from pg_class where relname='表名';

3)查询有哪些复制槽
select * from pg_replication_slots;

4)删除复制槽
select pg_drop_replication_slot('slot_name');  

5)排查wal日志缓冲区是否设置合理(buffer_full_ratio  在10%以下比较好)
查看wal日志缓冲区大小
show wal_buffers;

查看wal占比
show shared_buffers
SELECT
  -- 获取wal_buffers配置(块数→转换为MB)
  (s.setting::numeric * 8192) / (1024 * 1024) AS wal_buf_size_mb,
  -- 因缓冲区满触发的刷盘次数
  w.wal_buffers_full,
  -- 总WAL写入次数
  w.wal_write,
  -- 缓冲区满刷盘占比(核心判断指标)
  CASE
    WHEN w.wal_write > 0 THEN (w.wal_buffers_full::numeric / w.wal_write) * 100
    ELSE 0
  END AS buffer_full_ratio
FROM
  pg_settings s,
  pg_stat_wal w
WHERE
  s.name = 'wal_buffers';  -- 精准匹配参数名

6)排查共享数据缓冲区(backend_flush_ratio 在20%以下比较好)
查看共享缓冲区大小
show shared_buffers

查看共享占比
SELECT
  buffers_checkpoint,
  buffers_backend,
  buffers_clean,
  buffers_alloc,
  -- 计算用户同步刷盘占比(越高越差)
  round(buffers_backend::numeric / (buffers_checkpoint + buffers_backend + buffers_clean + 0.01) * 100, 2) AS backend_flush_ratio 
FROM pg_stat_bgwriter;

7)查看wal缓冲区写满次数
SELECT wal_buffers_full from pg_stat_wal;


8)排查flinkcdc同步延迟和日志堆积
 SELECT
  -- 基础标识字段
  prs.slot_name,        -- 复制槽名称
  psr.state,            -- 同步状态(catchup/streaming等)
  psr.client_addr,      -- CDC客户端/备库IP
  prs.active,           -- 复制槽是否活跃(true/false)
  -- 时间维度延迟:各阶段卡了多久(格式化到毫秒,空值补0)
  COALESCE(to_char(psr.write_lag, 'HH24:MI:SS.ms'), '00:00:00.000') AS write_lag_time,    -- 写盘延迟-时间
  COALESCE(to_char(psr.flush_lag, 'HH24:MI:SS.ms'), '00:00:00.000') AS flush_lag_time,  -- 刷盘延迟-时间
  COALESCE(to_char(psr.replay_lag, 'HH24:MI:SS.ms'), '00:00:00.000') AS replay_lag_time,-- 回放延迟-时间
  -- 数据维度积压:各阶段剩余未处理数据量(自动适配B/KB/MB/GB)
  pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), psr.write_lsn)) AS write_backlog_size,   -- 写盘积压-数据量
  pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), psr.flush_lsn)) AS flush_backlog_size, -- 刷盘积压-数据量
  pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), psr.replay_lsn)) AS replay_backlog_size,-- 回放积压-数据量
  -- 核心总滞后指标
  pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), prs.restart_lsn)) AS slot_total_backlog_size, -- 复制槽总滞后-数据量
  pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), COALESCE(confirmed_flush_lsn, restart_lsn))) AS unflush_backlog_size, -- 未刷盘滞后-数据量
  prs.restart_lsn,      --confirmed_flush_lsn- 复制槽重启LSN(追赶起始位点)
  prs.confirmed_flush_lsn -- 已确认刷盘LSN(CDC持久化位点)
FROM
  pg_stat_replication psr
LEFT JOIN
  pg_replication_slots prs ON psr.pid = prs.active_pid
WHERE
  prs.slot_name IS NOT NULL  -- 过滤无对应复制槽的无效连接
  -- and  prs.slot_name like ''


9)查看复制槽位点的变化
current_lsn:复制槽的起始wal位点
restart_lsn:复制槽的恢复wal位点
如果restart_lsn长时间不变化,代表没有新数据进来
max_slot_wal_keep_size(默认是不限制),如果DBA做了限制, 达到阈值会触发复制槽失效
SELECT 
  slot_name,
  pg_current_wal_lsn() AS current_lsn,
  restart_lsn,
  pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn) AS lsn_diff
FROM pg_replication_slots 
WHERE slot_name = '';

详细排查请参见这篇文章https://blog.csdn.net/ask_baidu/article/details/157063833

4.java关键代码

        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
//        env.setParallelism(1);
        // enable checkpoint
        env.enableCheckpointing(120000);
        env.getCheckpointConfig().setCheckpointTimeout(600000L);
        env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
        env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);
        env.setRestartStrategy(RestartStrategies.fixedDelayRestart(3, Time.milliseconds(10)));
        String tableList = String.join(",", tableNames.stream().map(c ->    SourceConfig.SCHEMA + "." + c).collect(Collectors.toList()));
        Properties proper = new Properties();
        proper.setProperty("snapshot.mode", "never");  // never initial
//        proper.setProperty("debezium.slot.name", slotName);
        proper.setProperty("slot.drop.on.stop", "true");
        proper.setProperty("database.serverTimezone", "GMT+8"); //设置时区
        proper.setProperty("decimal.handling.mode", "STRING");

        SourceFunction<String> sourceFunction = PostgreSQLSource.<String>builder()
                .hostname(SourceConfig.HOST_NAME)
                .port(SourceConfig.PORT)
                .database(SourceConfig.DATABASE_PRODUCE) // monitor postgresdatabase
                .schemaList(SourceConfig.SCHEMA) // monitor inventory schema
//                .tableList(tableList)
                .username(SourceConfig.USER_NAME)
                .password(SourceConfig.PASSWORD)
                .decodingPluginName("pgoutput") // pg解码插件
                .slotName(slotName + "_" + SourceConfig.DATABASE_PRODUCE) // 复制槽名称 不能重复(每个任务对应唯一的复制槽)
                .deserializer(new JsonDebeziumDeserializationSchema()) // converts SourceRecord to JSON String
                .debeziumProperties(proper).build();
                
           SingleOutputStreamOperator<String> process = cdcSource.process(new ProcessFunction<String, String>() {

            @Override
            public void processElement(String row, ProcessFunction<String, String>.Context context, Collector<String> collector) throws Exception {
			// row就是过来的日志
			}
           }).name("test" + "_" + SourceConfig.DATABASE_PRODUCET).setParallelism(15);
env.execute("test" + SourceConfig.DATABASE_PRODUCE);// job任务名称


5.pom.xml

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
         xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
         xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
    <modelVersion>4.0.0</modelVersion>

    <groupId>com.test</groupId>
    <artifactId>test-cdc</artifactId>
    <version>1.0-SNAPSHOT</version>

    <properties>
        <java.version>1.8</java.version>
        <maven.compiler.source>${java.version}</maven.compiler.source>
        <maven.compiler.target>${java.version}</maven.compiler.target>
        <fastjson.vsersion>2.0.52</fastjson.vsersion>
        <druid.version>1.2.15</druid.version>
        <flink.version>1.18.0</flink.version>
        <flinkcdc.vsersion>3.0.1</flinkcdc.vsersion>
        <scala.version>2.12</scala.version>
        <lombok.version>1.18.20</lombok.version>
        <postgresql.version>42.2.12</postgresql.version>
    </properties>

    <dependencies>
        <dependency>
            <groupId>org.postgresql</groupId>
            <artifactId>postgresql</artifactId>
            <version>${postgresql.version}</version>
        </dependency>
        <dependency>
            <groupId>com.ververica</groupId>
            <artifactId>flink-connector-postgres-cdc</artifactId>
            <version>${flinkcdc.vsersion}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-table-api-java-bridge</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-table-planner_${scala.version}</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-table-runtime</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-table-common</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-streaming-java</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-clients</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-cep</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-java</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.flink</groupId>
            <artifactId>flink-json</artifactId>
            <version>${flink.version}</version>
        </dependency>
        <dependency>
            <groupId>com.alibaba</groupId>
            <artifactId>fastjson</artifactId>
            <version>${fastjson.vsersion}</version>
        </dependency>
        <dependency>
            <groupId>org.projectlombok</groupId>
            <artifactId>lombok</artifactId>
            <version>${lombok.version}</version>
            <scope>provided</scope>
        </dependency>

        <dependency>
            <groupId>com.alibaba</groupId>
            <artifactId>druid-spring-boot-starter</artifactId>
            <version>${druid.version}</version>
        </dependency>
        <dependency>
            <groupId>org.apache.doris</groupId>
            <artifactId>flink-doris-connector-1.18</artifactId>
            <version>24.0.0</version>
        </dependency>
    </dependencies>

    <build>
        <plugins>
            <plugin>
                <groupId>org.apache.maven.plugins</groupId>
                <artifactId>maven-assembly-plugin</artifactId>
                <version>3.0.0</version>
                <configuration>
                    <descriptorRefs>
                        <descriptorRef>jar-with-dependencies</descriptorRef>
                    </descriptorRefs>
                </configuration>
                <executions>
                    <execution>
                        <id>make-assembly</id>
                        <phase>package</phase>
                        <goals>
                            <goal>single</goal>
                        </goals>
                    </execution>
                </executions>
            </plugin>
        </plugins>
    </build>
</project>

6. 问题总结
1)日志堆积导致max_slot_wal_keep_size超出阈值触发连锁反应
checkpoint设置太大->遇到PG峰值(瞬时大量日志)->max_slot_wal_keep_size超出阈值后->断开Flink CDC的连接->Flink CDC复制槽失效->Flink子节点挂掉->solt不足,所有子节点挂掉
Flink CDC报错日志(表象):ERROR io.debezium.pipeline.ErrorHandler [] - Producer failure
org.postgresql.util.PSQLException: Database connection failed when reading from copy
Caused by: java.io.EOFException
INFO Connection gracefully closed
PG日志(表象):LOG: invalidating slot “Flink CDC复制槽名称” because its restart_lsn ABF1/2FDC958 exceeds max_slot_wal_keep_size
Flink 任务异常界面(表象):Caused by: org.postgresql.util.PSQLException: FATAL: terminating connection due to administrator command

(本质)org.apache.kafka.connect.errors.ConnectException: Unable to obtain valid replication slot. Make sure there are no long-running transactions running in parallel as they may hinder the allocation of the replication slot when starting this connector(PostgresConnection)
原因:次因:max_slot_wal_keep_size默认是没有限制的,如果DBA做了限制,超过阈值会终止Flink CDC的连接
举例:max_slot_wal_keep_size = 50GB
slot_total_backlog_size = 45GB ( 复制槽总滞后大小,slot_total_backlog_size 参考本节上面sql)
replay_backlog_size = 30GB ( 回放积压大小,replay_backlog_size 参考本节上面sql)
说明:Flink CDC本身落后 30GB,PG内部还有15GB WAL没发给它
总落后 45GB,再涨 5GB,就会达到max_slot_wal_keep_size 阈值
主因:1.checkpoint设置过大
2.长事务(小时级别以上的)
解决:1.调大max_slot_wal_keep_size
2.提高checkpoint的频率(值设小点)
3.避免长事务阻塞位点变化导致日志持续堆积,最好分批提交
总结:1.复制槽是Flink CDC和PG建立复制连接的核心关键入口
2.checkpoint是CDC重磅角色
2)java Metaspace OOM
报错日志:com.alibaba.fastjson.JSONException: create objectReader error, objectType
Caused by: java.lang.OutOfMemoryError: Metaspace. The metaspace out-of-memory error has occurred. This can mean two things: either the job requires a larger size of JVM metaspace to load classes or there is a class loading leak. In the first case ‘taskmanager.memory.jvm-metaspace.size’ configuration option should be increased. If the error persists (usually in cluster after several job (re-)submissions) then there is probably a class loading leak in user code or some of its dependencies which has to be investigated and fixed. The task executor has to be shutdown…
原因:java元空间不足
解决:flink-config.xml加入taskmanager.memory.jvm-metaspace.size: 1024m(自己根据实际情况调整)

官方配置参数:
https://docs.confluent.io/kafka-connectors/debezium-postgres-source/current/postgres_source_connector_config.html?use_xbridge3=true&loader_name=forest&need_sec_link=1&sec_link_scene=im&theme=light#required-properties

更多推荐