flinkcdc抽取postgres数据
摘要
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
更多推荐
所有评论(0)