![]() | 博主历时三年精心创作的《大数据平台架构与原型实现:数据中台建设实战》一书现已由知名IT图书品牌电子工业出版社博文视点出版发行,点击《重磅推荐:建大数据平台太难了!给我发个工程原型吧!》了解图书详情,京东购书链接:https://item.jd.com/12677623.html,扫描左侧二维码进入京东手机购书页面。 |
本文介绍的整体方案选型是:使用 Kafka Connect 的 Debezium MySQL Source Connector 将 MySQL 的 CDC 数据 (Avro 格式)接入到 Kafka 之后,通过 Flink 读取并解析这些 CDC 数据,其中,数据是以 Confluent 的 Avro 格式存储的,也就是说,Avro 格式的数据在写入到 Kafka 以及从 Kafka 读取时,都需要和 Confluent Schema Registry 进行交互,从而获取 Schema 信息,消息经 Flink 读取后会写入到 Hudi 表,从而完成全部的数据接入工作。
1. 前置依赖
本文不会展开介绍 CDC 数据进入 Kafka 之前的操作,此部分可以参考 [《CDC一键入湖:当 Apache Hudi DeltaStreamer 遇见 Serverless Spark》](https://laurence.blog.csdn.net/article/details/132011197) 一文的前半部分架构以及第 2 节环境准备部分的介绍,以下是前半部分数据管道使用到的相关组件的构建方法和文档:
①MySQL:如果仅以测试为目的,建议使用Debezium提供的官方Docker镜像,构建操作可参考其官方文档(下文将给出的操作示例所处理的CDC数据就是自于该MySQL镜像中的inventory数据库);
②Kafka Connect:如果仅以测试为目的,建议使用Confluent提供的官方Docker镜像,构建操作可参考其官方文档,或者使用AWS上托管的Kafka Connct:Amazon MSK Connect。需要提醒的是:Kafka Connect上必须安装Debezium MySQL Connector和Confluent Avro Converter两个插件,因此需要在官方镜像的基础上手动添加这两个插件;
③Confluent Schema Registry:如果仅以测试为目的,建议使用Confluent提供的官方Docker镜像,构建操作可参考其官方文档;
④Kafka:如果仅以测试为目的,建议使用Confluent提供的官方Docker镜像,构建操作可参考其官方文档,或者使用AWS上托管的Kafka:Amazon MSK
本文讨论的核心是:Flink 如何从 Kafka 中读取并解析 Debezium Confluent Avro 消息并写入到 Hudi 表中。
2. 环境准备
为完成既定目标,我们将会使用和依赖到多个 Flink 插件,包括:Flink Sql Kafka Connector、Flink Hudi Connector、Flink Hive Connector、Flink ‘debezium-avro-confluent’ Format Support,以下脚本将会安装这些插件及其依赖包:
# install flink kafka connector for flink sql client
# only run on master node is enough, owner of flink home dir is 'flink' user
sudo -u flink wget https://repo.maven.apache.org/maven2/org/apache/flink/flink-sql-connector-kafka/1.17.1/flink-sql-connector-kafka-1.17.1.jar -P /usr/lib/flink/lib/# install flink hudi connector for flink sql client
# only run on master node is enough, owner of flink home dir is 'flink' user
sudo -u flink wget https://repo1.maven.org/maven2/org/apache/hudi/hudi-flink1.17-bundle/0.14.0/hudi-flink1.17-bundle-0.14.0.jar -P /usr/lib/flink/lib/# install flink hudi connector for flink sql client
# only run on master node is enough, owner of flink home dir is 'flink' user
# refer to this doc: https://docs.aws.amazon.com/emr/latest/ReleaseGuide/flink-configure.html
sudo -u flink cp /usr/lib/hive/lib/antlr-runtime-3.5.2.jar /usr/lib/flink/lib
sudo -u flink cp /usr/lib/hive/lib/hive-exec-3.1.3*.jar /usr/lib/flink/lib
sudo -u flink cp /usr/lib/hive/lib/libfb303-0.9.3.jar /usr/lib/flink/lib
sudo -u flink cp /usr/lib/flink/opt/flink-connector-hive_2.12-1.17.1-amzn-1.jar /usr/lib/flink/lib# install flink 'debezium-avro-confluent' format for flink sql client
# only run on master node is enough, owner of flink home dir is 'flink' user
# refer to this doc: https://blog.csdn.net/bluishglc/article/details/135863249 , section 3.2
sudo -u flink unzip jar_files.zip -d /usr/lib/flink/lib/
接下来,清空 Hudi 表目标位置上的文件,停止正在运行中的 Yarn App,启动一个新的 Flink Yarn Session:
echo "clean hudi table target location..."
aws s3 rm --recursive s3://glc-flink-hudi-test/sink_hudi_orders
echo "Kill all running apps..."
for appId in $(yarn application -list -appStates RUNNING 2>1 | awk 'NR > 2 { print $1 }'); doyarn application -kill $appId &> /dev/null
done
echo "start a flink yarn session..."
flink-yarn-session -d
再启动一个 Flink SQL Client,准备输入 SQL:
/usr/lib/flink/bin/sql-client.sh embedded shell
3. Flink 从 Kafka 中读取并解析 Debezium Confluent Avro 消息
关于此环节的配置和操作,实际上我们已经在《Flink 集成 Debezium Confluent Avro ( format=debezium-avro-confluent )》一文给出了详细的介绍,所以请移步此文了解详细操作。但是,我们稍微修改了一下SQL,启用了 Hive MetaStore,也重构了一下表名。在执行实际操作前,需先安装 Flink Hive Connector,此操作请参考 《Flink 集成和使用 Hive Metastore》。以下是《Flink 集成 Debezium Confluent Avro ( format=debezium-avro-confluent )》一文的 SQL 基础之上,引入 Hive Metastore 并重构了表名后的版本:
CREATE CATALOG hive WITH ('type' = 'hive','hive-conf-dir' = '/etc/hive/conf'
);USE CATALOG hive;DROP TABLE IF EXISTS src_cdc_orders;
CREATE TABLE IF NOT EXISTS src_cdc_orders (-- a few columns mapped to the Avro fields of the Kafka valueorder_number int,order_date int,purchaser int,quantity int,product_id int,event_time TIMESTAMP(3) METADATA FROM 'timestamp' VIRTUAL
) WITH ('connector' = 'kafka','topic' = 'osci.mysql-server-3.inventory.orders','properties.bootstrap.servers' = 'b-2.oscimskcluster1.cedsl9.c20.kafka.us-east-1.amazonaws.com:9092,b-3.oscimskcluster1.cedsl9.c20.kafka.us-east-1.amazonaws.com:9092,b-1.oscimskcluster1.cedsl9.c20.kafka.us-east-1.amazonaws.com:9092','properties.group.id' = 'flink','scan.startup.mode' = 'earliest-offset','format' = 'debezium-avro-confluent','debezium-avro-confluent.schema-registry.url' = 'http://10.0.13.30:8085'
);
4. Flink 流式读取 Kafka 中的数据写入 Hudi 表
当数据写入到 Kafka 后,接下来就是创建一个 Hudi 表,并将上述 `src_cdc_orders` 表数据写入到 Huidi 表就可以了:
DROP TABLE IF EXISTS sink_hudi_orders;CREATE TABLE IF NOT EXISTS sink_hudi_orders (order_number int PRIMARY KEY NOT ENFORCED,order_date int,purchaser int,quantity int,product_id int,event_time TIMESTAMP(3)
) WITH ('connector' = 'hudi','path' = 's3://glc-flink-hudi-test/sink_hudi_orders','table.type' = 'MERGE_ON_READ','precombine.field' = 'event_time','changelog.enabled' = 'true'
);insert into sink_hudi_orders select * from src_cdc_orders;select * from sink_hudi_orders/*+ OPTIONS('read.streaming.enabled'='true', 'read.streaming.check-interval' = '2', 'read.streaming.skip_compaction' = 'true', 'read.streaming.start-commit' = 'earliest')*/;
在上面的 SQL 中,开头部分的几个 SET 语句非常重要,它们用于设置 Checkpoint,如果没有这些设置,写入操作都不会提交,看到的状况就是:作业流一直运行,没有报错,但是 Hudi 表不会有任何数据,关于这一问题以及这些 SET 语句的解释,已在《Flink 读取 Kafka 消息写入 Hudi 表无报错但没有写入任何记录的解决方法》一文做了详细介绍,请移步此文了解更多细节。
下面演示了使用 Flink SQL 创建 Kafka 源表,然后读取 Kafka 中的 CDC 数据并流式写入到 Hudi 表的全过程:

5. 常见错误
1. 作业没有报错,但 Hudi 表中无数据(记录不会写入到 Hudi 表中)
该问题非常典型,也几乎总会遇到,详细解释和处理方法已经总结在 《Flink 读取 Kafka 消息写入 Hudi 表无报错但没有写入任何记录的解决方法》一文中。
2. MOR 表,开启 ‘changelog.enabled’ = ‘true’ ,无 -D 删除记录
关于这一问题,已详细记录在 《Flink 流式读取 Debezium CDC 数据写入 Hudi 表无法处理 -D / Delete 消息》 中,目前暂无解决方案,推测可能是 Flink 尚未实现该功能,或者还有什么关键配置没有配对。
