ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Flink CDC 实战指南:从零搭建 PostgreSQL 到 Fluss 的全库实时同步与 Schema 演进流水线

Flink CDC 实战指南:从零搭建 PostgreSQL 到 Fluss 的全库实时同步与 Schema 演进流水线 Flink CDC 实战指南从零搭建 PostgreSQL 到 Fluss 的全库实时同步与 Schema 演进流水线【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc本文以 Flink CDC 官方教程为主线完整演示如何基于 Flink 2.2.0 Standalone 集群与 Flink CDC CLI构建一条从 PostgreSQL 到 Fluss 的 Streaming ELT 流水线覆盖全库同步、SASL 安全接入以及增列/删列/改名列三类 Schema 变更的自动演进。读完后你可以不写一行 Java/Scala 代码仅凭 SQL 和 YAML 配置文件即可完成建库表数据 → 提交同步任务 → 在 Fluss 中查询 → 验证表结构同步的全流程并能从源码层面理解schema-change.enabled、schema.change.behavior: LENIENT等关键参数背后的实现机制。一、环境与准备整个教程运行在一台安装好 Docker 的 Linux 或 macOS 机器上需要准备三块环境Flink 2.2.0 Standalone 集群、Docker Compose 拉起的 Fluss 集群coordinator-server tablet-server ZooKeeper以及一个开启逻辑复制的 PostgreSQL 14.5 源库。1.1 准备 Flink Standalone 集群下载 Flink 2.2.0 发行包可从 Apache 官方归档站点获取flink-2.2.0-bin-scala_2.12.tgz解压得到flink-2.2.0目录并进入该目录cd flink-2.2.0编辑conf/config.yaml追加以下配置以开启 3 秒一次的 checkpointFlink CDC 全库同步依赖 checkpoint 保证 exactly-once 语义与 LSN 提交execution: checkpointing: interval: 3s启动集群./bin/start-cluster.sh启动成功后可访问http://localhost:8081/查看 Flink Web UI。重复执行start-cluster.sh可以增加 TaskManager 数量。1.2 用 Docker Compose 准备 Fluss 与 PostgreSQL创建docker-compose.yml内容如下。其中 Fluss 侧使用了apache/fluss:0.9.0-incubating镜像并开启了 SASL/PLAIN 认证developer用户PostgreSQL 侧通过启动参数开启了逻辑复制所需的wal_levellogical、max_replication_slots5、max_wal_senders5services: # Fluss cluster coordinator-server: image: apache/fluss:0.9.0-incubating command: coordinatorServer depends_on: - zookeeper environment: - | FLUSS_PROPERTIES zookeeper.address: zookeeper:2181 bind.listeners: INTERNAL://coordinator-server:0, CLIENT://coordinator-server:9123 advertised.listeners: CLIENT://localhost:9123 internal.listener.name: INTERNAL remote.data.dir: /tmp/fluss/remote-data security.protocol.map: CLIENT:SASL, INTERNAL:PLAINTEXT security.sasl.enabled.mechanisms: PLAIN security.sasl.plain.jaas.config: org.apache.fluss.security.auth.sasl.plain.PlainLoginModule required user_adminadmin-pass user_developerdeveloper-pass ; super.users: User:admin ports: - 9123:9123 tablet-server: image: apache/fluss:0.9.0-incubating command: tabletServer depends_on: - coordinator-server environment: - | FLUSS_PROPERTIES zookeeper.address: zookeeper:2181 bind.listeners: INTERNAL://tablet-server:0, CLIENT://tablet-server:9123 advertised.listeners: CLIENT://localhost:9124 internal.listener.name: INTERNAL tablet-server.id: 0 kv.snapshot.interval: 0s data.dir: /tmp/fluss/data remote.data.dir: /tmp/fluss/remote-data security.protocol.map: CLIENT:SASL, INTERNAL:PLAINTEXT security.sasl.enabled.mechanisms: PLAIN security.sasl.plain.jaas.config: org.apache.fluss.security.auth.sasl.plain.PlainLoginModule required user_adminadmin-pass user_developerdeveloper-pass ; super.users: User:admin ports: - 9124:9123 zookeeper: restart: always image: zookeeper:3.9.2 # PostgreSQL postgres: image: postgres:14.5 environment: POSTGRES_USER: root POSTGRES_PASSWORD: password POSTGRES_DB: postgres ports: - 5432:5432 volumes: - postgres_data:/var/lib/postgresql/data command: - postgres - -c - wal_levellogical - -c - max_replication_slots5 - -c - max_wal_senders5 - -c - hot_standbyon volumes: postgres_data:各服务的作用Flusscoordinator-server、tablet-server、zookeeper目标数据湖仓存储。coordinator-server 负责元数据管理tablet-server 负责数据存储客户端通过localhost:9123coordinator连接。PostgreSQL源数据库wal_levellogical是 CDC 逻辑解码的硬性前提。在docker-compose.yml所在目录执行docker-compose up -d用docker ps确认所有容器运行正常。1.3 初始化 PostgreSQL 源数据连接 PostgreSQL密码为passwordpsql -h localhost -p 5432 -U root postgres创建adb数据库并切换过去CREATE DATABASE adb; \c adb创建两个 schemahr、sales及对应表并插入初始数据-- Create schemas CREATE SCHEMA hr; CREATE SCHEMA sales; -- Create tables CREATE TABLE hr.employees( ID INT PRIMARY KEY NOT NULL, NAME TEXT NOT NULL, AGE INT NOT NULL, ADDRESS CHAR(50), SALARY REAL ); CREATE TABLE sales.orders( ID INT PRIMARY KEY NOT NULL, PRODUCT TEXT NOT NULL, QUANTITY INT NOT NULL, REGION CHAR(50), AMOUNT REAL ); -- Insert data INSERT INTO hr.employees (ID, NAME, AGE, ADDRESS, SALARY) VALUES (1, Paul, 32, California, 20000.00); INSERT INTO sales.orders (ID, PRODUCT, QUANTITY, REGION, AMOUNT) VALUES (1, Laptop, 5, East, 49999.50);这里刻意构造了两个 schema 下的两张表用于验证全库同步能力——后续任务会把adb库下所有 schema 的所有表都同步到 Fluss。二、编写并提交 Flink CDC 同步任务2.1 下载并放置组件包下载 Flink CDC 二进制发行包flink-cdc-version-bin.tar.gz从 Apache 官方下载渠道获取解压得到flink-cdc-version目录其中包含bin、lib、log、conf四个子目录。下载两个 Pipeline Connector 包稳定版可从 Maven 仓库获取SNAPSHOT 版本需自行基于 master 或 release 分支构建并放入Flink CDC Home 的lib目录注意不是 Flink Home 的lib目录flink-cdc-pipeline-connector-postgresflink-cdc-pipeline-connector-fluss这两个 Connector 在源码仓库中分别对应flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-postgres/与flink-cdc-connect/flink-cdc-pipeline-connectors/flink-cdc-pipeline-connector-fluss/两个模块。2.2 编写 pipeline 配置 yaml下面是一份完整的postgres-to-fluss.yaml用于将adb全库同步到 Fluss################################################################################ # Description: Sync Postgres all tables to Fluss ################################################################################ source: type: postgres hostname: localhost port: 5432 username: root password: password tables: adb.\.*.\.* decoding.plugin.name: pgoutput slot.name: pgtest schema-change.enabled: true sink: type: fluss bootstrap.servers: localhost:9123 properties.client.security.protocol: sasl properties.client.security.sasl.mechanism: PLAIN properties.client.security.sasl.username: developer properties.client.security.sasl.password: developer-pass pipeline: name: Postgres to Fluss Pipeline parallelism: 2 schema.change.behavior: LENIENT关键参数说明source 部分tables: adb.\.*.\.*用正则匹配adb库下所有 schema 的所有表实现全库同步。注意 CDC 只支持连接单一数据库所有表必须属于同一个库.是库/schema/表三级分隔符正则中需要匹配任意字符时必须写成\.转义。decoding.plugin.name: pgoutputPostgreSQL 逻辑解码插件取值为decoderbufs或pgoutput。本教程必须用pgoutput因为源端 Schema 变更推断依赖 pgoutput 的 Relation 消息。slot.name: pgtest逻辑复制槽名称必须小写字母、数字、下划线服务端通过该槽向 Connector 推送变更事件。schema-change.enabled: true开启源端 Schema 变更推断。该选项定义在 PostgresDataSourceOptions.java 中默认值为false。开启后Connector 会把 pgoutput 的 Relation 消息与缓存 schema 做比对推断出加列、删列、改列名、改列类型等事件。完整的 Postgres 源端参数含scan.startup.mode、scan.incremental.snapshot.chunk.size等可参考 postgres.md。sink 部分Fluss Sink 的内置参数定义在 FlussDataSinkOptions.java核心项包括bootstrap.servers必填Fluss 集群引导地址本教程为localhost:9123bucket.key可选按表指定分桶键格式database1.table1:key1,key2;database1.table2:key3分桶键必须是主键表的非分区主键子集未指定时主键表默认以主键去掉分区键作为分桶键bucket.num可选按表指定桶数格式database1.table1:4;database1.table2:8properties.table.*/properties.client.*分别透传 Fluss 表选项与客户端选项。本教程正是通过properties.client.security.*四个键配置 SASL/PLAIN 认证与 compose 文件中user_developerdeveloper-pass的用户对应。这些配置在工厂类中经 FlussConfigUtils.java 解析成表标识 → 分桶键列表/桶数的映射最终注入 FlussDataSink.java。从源码结构看FlussDataSink同时提供了三部分能力数据写入FlussSink底层走 Fluss Java Client、Schema 变更应用FlussMetaDataApplier以及按分桶键哈希的数据分片函数FlussHashFunctionProvider。pipeline 部分name/parallelism任务名与并行度。schema.change.behavior: LENIENT必须显式声明。这一点容易被忽略原因是 Flink CDC 发行包自带的默认配置 flink-cdc.yaml 中写有# Parallelism of the pipeline parallelism: 4 # Behavior for handling schema change events from source schema.change.behavior: EVOLVE如果任务 yaml 不显式指定任务级schema.change.behavior可能被默认conf.yaml中的EVOLVE覆盖而 EVOLVE 模式下 Fluss 侧对删列、改列名等事件的处理能力有限同步会失败。五种行为模式exception/evolve/try_evolve/lenient/ignore的完整语义见 schema-evolution.md。选择 LENIENT 模式后Fluss 侧各 Schema 事件的实际处理方式依据 fluss.md 的 Usage Notes加列新列追加到 Fluss 表末尾删列不物理删除drop 操作被忽略后续写入将该列置为 null改列名等价转换为新增新列 旧列改为可空两个事件改列类型不支持。这与 FlussMetaDataApplier.java 的实现一致它只处理CreateTableEvent、DropTableEvent、AddColumnEvent三类事件且加列时强制校验位置必须为LAST否则抛出异常并提示使用 LENIENT 模式。建表时还会做sanityCheck比对推断出的主键、分桶键、分区键与 Fluss 中已有表是否一致防止意外的 Schema 漂移。2.3 提交任务在 Flink CDC 目录下用 CLI 提交任务bash bin/flink-cdc.sh postgres-to-fluss.yaml提交成功后返回Pipeline has been submitted to cluster. Job ID: ae30f4580f1918bebf16752d4963dc54 Job Description: Postgres to Fluss Pipeline此时在 Flink Web UIhttp://localhost:8081/中可以找到名为Postgres to Fluss Pipeline的运行中作业。全库同步流程会先对hr.employees与sales.orders做增量快照读取然后切换到复制槽pgtest的增量日志读取。三、在 Fluss 中查询同步结果查询侧需要一个能连接 Fluss 的 Flink SQL Client将 Fluss 的 Flink 2.2 兼容连接器 jarfluss-flink-2.2-0.9.0-incubating.jar可从 Maven 中央仓库获取放入 Flink 的lib目录。启动 SQL Clientbin/sql-client.sh创建 Fluss Catalog使用与 compose 文件一致的developer用户认证并列出数据库SET execution.runtime-mode batch; SET sql-client.execution.result-mode tableau; CREATE CATALOG developer_catalog WITH ( type fluss, bootstrap.servers localhost:9123, client.security.protocol SASL, client.security.sasl.mechanism PLAIN, client.security.sasl.username developer, client.security.sasl.password developer-pass ); USE CATALOG developer_catalog; SHOW DATABASES; --------------- | database name | --------------- | fluss | | hr | | sales | ---------------hr与sales两个 database 正是源端 PostgreSQL 中adb库的两个 schema 自动映射而来——Fluss 侧在收到CreateTableEvent时会先createDatabase不存在则创建再建表。查询同步过来的表SELECT * FROM developer_catalog.hr.employees LIMIT 20; -------------------------------------------------------- | id | name | age | address | salary | -------------------------------------------------------- | 1 | Paul | 32 | California ... | 20000.0 | --------------------------------------------------------四、Schema 变更同步实战PostgreSQL 的 Schema 变更捕获是数据驱动的仅执行 DDL 并不会立即产生事件必须等到下一条 DML 消息触发 pgoutput 插件发出新的 Relation 消息Connector 才能对比出 Schema 差异。因此以下每个实验都是先 ALTER 表、再 INSERT 数据的组合。连接源库psql -h localhost -p 5432 -U root adb4.1 添加列ALTER TABLE hr.employees ADD COLUMN EMAIL TEXT, ADD COLUMN DEPARTMENT TEXT; INSERT INTO hr.employees (ID, NAME, AGE, ADDRESS, SALARY, EMAIL, DEPARTMENT) VALUES (4, David, 32, Guangzhou, 8000.0, davidexample.com, IT), (5, Eva, 27, Hangzhou, 7100.0, evaexample.com, HR);在 Fluss 中查询新列已自动追加在表末尾存量行在新列上显示NULLSELECT * FROM developer_catalog.hr.employees LIMIT 20; ---------------------------------------------------------------------------------------- | id | name | age | address | salary | email | department | ---------------------------------------------------------------------------------------- | 1 | Paul | 32 | California ... | 20000.0 | NULL | NULL | | 4 | David | 32 | Guangzhou ... | 8000.0 | davidexample.com | IT | | 5 | Eva | 27 | Hangzhou ... | 7100.0 | evaexample.com | HR | ----------------------------------------------------------------------------------------对应链路pgoutput Relation 消息 → 源端推断出AddColumnEvent→ LENIENT 模式放行 →FlussMetaDataApplier.applyAddColumnTable调用 Fluss Admin 的alterTable以ColumnPosition.last()追加列。4.2 删除列ALTER TABLE hr.employees DROP COLUMN ADDRESS; INSERT INTO hr.employees (ID, NAME, AGE, SALARY, EMAIL, DEPARTMENT) VALUES (6, Frank, 35, 9000.0, frankexample.com, Finance), (7, Grace, 29, 7600.0, graceexample.com, Marketing);查询 Fluss可以看到 LENIENT 模式下旧数据中address值被保留、新写入的行该列为NULL列本身未被物理删除SELECT * FROM developer_catalog.hr.employees LIMIT 20; ---------------------------------------------------------------------------------------- | id | name | age | address | salary | email | department | ---------------------------------------------------------------------------------------- | 1 | Paul | 32 | California ... | 20000.0 | NULL | NULL | | 4 | David | 32 | Guangzhou ... | 8000.0 | davidexample.com | IT | | 5 | Eva | 27 | Hangzhou ... | 7100.0 | evaexample.com | HR | | 6 | Frank | 35 | NULL | 9000.0 | frankexample.com | Finance | | 7 | Grace | 29 | NULL | 7600.0 | graceexample.com | Marketing | ----------------------------------------------------------------------------------------4.3 重命名列ALTER TABLE hr.employees RENAME COLUMN EMAIL TO WORK_EMAIL; INSERT INTO hr.employees (ID, NAME, AGE, SALARY, WORK_EMAIL, DEPARTMENT) VALUES (8, Henry, 31, 8800.0, henryexample.com, Sales), (9, Ivy, 26, 6900.0, ivyexample.com, Support);查询 Fluss可以看到 LENIENT 模式把改列名转换为新增work_email列 旧email列保留的组合效果SELECT * FROM developer_catalog.hr.employees LIMIT 20; ----------------------------------------------------------------------------------------------------------- | id | name | age | address | salary | email | department | work_email | ----------------------------------------------------------------------------------------------------------- | 1 | Paul | 32 | California ... | 20000.0 | NULL | NULL | NULL | | 4 | David | 32 | Guangzhou ... | 8000.0 | davidexample.com | IT | NULL | | 5 | Eva | 27 | Hangzhou ... | 7100.0 | evaexample.com | HR | NULL | | 6 | Frank | 35 | NULL | 9000.0 | frankexample.com | Finance | NULL | | 7 | Grace | 29 | NULL | 7600.0 | graceexample.com | Marketing | NULL | | 8 | Henry | 31 | NULL | 8800.0 | NULL | Sales | henryexample.com | | 9 | Ivy | 26 | NULL | 6900.0 | NULL | Support | ivyexample.com | -----------------------------------------------------------------------------------------------------------如果希望更细粒度地控制哪些 Schema 事件需要同步例如只允许加列、禁止删列还可以在sink块使用include.schema.changes/exclude.schema.changes按事件类型过滤详见 schema-evolution.md 的 Per-Event Type Control 一节。五、清理环境教程完成后在docker-compose.yml所在目录停止并清理所有容器docker-compose down -v在 Flinkflink-2.2.0目录停止集群./bin/stop-cluster.sh六、小结本教程完成了一条可复现的 PostgreSQL → Fluss 全库同步链路Flink 2.2.0 Standalone 集群提供执行与 checkpoint 能力Docker Compose 提供带 SASL 认证的 Fluss 目标端与开启逻辑复制的 PostgreSQL 源端Flink CDC CLI 加一份 YAML 即可完成全库同步与三类列级 Schema 变更的自动演进。实践中需要记住的三个要点源端必须decoding.plugin.name: pgoutput且显式打开schema-change.enabled任务级schema.change.behavior: LENIENT必须显式声明避免被 flink-cdc.yaml 的默认EVOLVE覆盖Fluss 侧 LENIENT 语义是加列末尾追加、删列只置空不删除、改列名转为新增列改列类型不受支持。类型映射的完整对照如VARCHAR → STRING、VARBINARY → BYTES可参考 fluss.md 的 Data Type Mapping 一节。【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表