Apache Doris 整合 FLINK CDC + Iceberg 构建实时湖仓一体的联邦查询

1概况

本文展示如何使用 Flink CDC + Iceberg + Doris 构建实时湖仓一体的联邦查询分析,Doris 1.1版本提供了Iceberg的支持,本文主要展示Doris和Iceberg怎么使用,大家按照步骤可以一步步完成。完整体验整个搭建操作的过程。

2系统架构

我们整理架构图如下,

1.首先我们从Mysql数据中使用Flink 通过 Binlog完成数据的实时采集

2.然后再Flink 中创建 Iceberg 表,Iceberg的元数据保存在hive里

3.最后我们在Doris中创建Iceberg外表

4.在通过Doris 统一查询入口完成对Iceberg里的数据进行查询分析,供前端应用调用,这里iceberg外表的数据可以和Doris内部数据或者Doris其他外部数据源的数据进行关联查询分析

Doris湖仓一体的联邦查询架构如下:

1.Doris 通过 ODBC 方式支持:MySQL,Postgresql,Oracle ,SQLServer

2.同时支持 Elasticsearch 外表

3.1.0版本支持Hive外表

4.1.1版本支持Iceberg外表

5.1.2版本支持Hudi 外表

3 创建MySQL数据库表并初始化数据

1CREATE DATABASE demo; 2USE demo; 3CREATE TABLE userinfo ( 4 id int NOT NULL AUTO_INCREMENT, 5 name VARCHAR(255) NOT NULL DEFAULT 'flink', 6 address VARCHAR(1024), 7 phone_number VARCHAR(512), 8 email VARCHAR(255), 9 PRIMARY KEY (`id`) 10)ENGINE=InnoDB ; 11INSERT INTO userinfo VALUES (10001,'user_110','Shanghai','13347420870', NULL); 12INSERT INTO userinfo VALUES (10002,'user_111','xian','13347420870', NULL); 13INSERT INTO userinfo VALUES (10003,'user_112','beijing','13347420870', NULL); 14INSERT INTO userinfo VALUES (10004,'user_113','shenzheng','13347420870', NULL); 15INSERT INTO userinfo VALUES (10005,'user_114','hangzhou','13347420870', NULL); 16INSERT INTO userinfo VALUES (10006,'user_115','guizhou','13347420870', NULL); 17INSERT INTO userinfo VALUES (10007,'user_116','chengdu','13347420870', NULL); 18INSERT INTO userinfo VALUES (10008,'user_117','guangzhou','13347420870', NULL); 19INSERT INTO userinfo VALUES (10009,'user_118','xian','13347420870', NULL);

4 创建Iceberg Catalog

1CREATE CATALOG hive_catalog WITH ( 2 'type'='iceberg', 3 'catalog-type'='hive', 4 'uri'='thrift://localhost:9083', 5 'clients'='5', 6 'property-version'='1', 7 'warehouse'='hdfs://localhost:8020/user/hive/warehouse' 8);

5 创建 Mysql CDC 表

1CREATE TABLE user_source ( 2 database_name STRING METADATA VIRTUAL, 3 table_name STRING METADATA VIRTUAL, 4 `id` DECIMAL(20, 0) NOT NULL, 5 name STRING, 6 address STRING, 7 phone_number STRING, 8 email STRING, 9 PRIMARY KEY (`id`) NOT ENFORCED 10 ) WITH ( 11 'connector' = 'mysql-cdc', 12 'hostname' = 'localhost', 13 'port' = '3306', 14 'username' = 'root', 15 'password' = 'MyNewPass4!', 16 'database-name' = 'demo', 17 'table-name' = 'userinfo' 18 );

6 创建Iceberg表

1---查看catalog 2show catalogs; 3---使用catalog 4use catalog hive_catalog; 5--创建数据库 6CREATE DATABASE iceberg_hive; 7--使用数据库 8use iceberg_hive; 9

7 创建表

1CREATE TABLE all_users_info ( 2 database_name STRING, 3 table_name STRING, 4 `id` DECIMAL(20, 0) NOT NULL, 5 name STRING, 6 address STRING, 7 phone_number STRING, 8 email STRING, 9 PRIMARY KEY (database_name, table_name, `id`) NOT ENFORCED 10 ) WITH ( 11 'catalog-type'='hive' 12 );

从CDC表里插入数据到Iceberg表里

1use catalog default_catalog; 23insert into hive_catalog.iceberg_hive.all_users_info select * from user_source;

我们去查询iceberg表

select * from hive_catalog.iceberg_hive.all_users_info

8 Doris 查询 Iceberg

8.1 创建Iceberg外表

1CREATE TABLE `all_users_info` 2ENGINE = ICEBERG 3PROPERTIES ( 4"iceberg.database" = "iceberg_hive", 5"iceberg.table" = "all_users_info", 6"iceberg.hive.metastore.uris" = "thrift://localhost:9083", 7"iceberg.catalog.type" = "HIVE_CATALOG" 8); 9 10

参数说明

•ENGINE 需要指定为 ICEBERG

•PROPERTIES 属性:

iceberg.hive.metastore.uris:Hive Metastore 服务地址

iceberg.database:挂载 Iceberg 对应的数据库名

iceberg.table:挂载 Iceberg 对应的表名,挂载 Iceberg database 时无需指定。

iceberg.catalog.type:Iceberg 中使用的 catalog 方式,默认为 HIVE_CATALOG,当前仅支持该方式,后续会支持更多的 Iceberg catalog 接入方式。

1mysql> CREATE TABLE `all_users_info` 2 -> ENGINE = ICEBERG 3 -> PROPERTIES ( 4 -> "iceberg.database" = "iceberg_hive", 5 -> "iceberg.table" = "all_users_info", 6 -> "iceberg.hive.metastore.uris" = "thrift://localhost:9083", 7 -> "iceberg.catalog.type" = "HIVE_CATALOG" 8 -> ); 9Query OK, 0 rows affected (0.23 sec) 1011mysql> select * from all_users_info; 12+---------------+------------+-------+----------+-----------+--------------+-------+ 13| database_name | table_name | id | name | address | phone_number | email | 14+---------------+------------+-------+----------+-----------+--------------+-------+ 15| demo | userinfo | 10004 | user_113 | shenzheng | 13347420870 | NULL | 16| demo | userinfo | 10005 | user_114 | hangzhou | 13347420870 | NULL | 17| demo | userinfo | 10002 | user_111 | xian | 13347420870 | NULL | 18| demo | userinfo | 10003 | user_112 | beijing | 13347420870 | NULL | 19| demo | userinfo | 10001 | user_110 | Shanghai | 13347420870 | NULL | 20| demo | userinfo | 10008 | user_117 | guangzhou | 13347420870 | NULL | 21| demo | userinfo | 10009 | user_118 | xian | 13347420870 | NULL | 22| demo | userinfo | 10006 | user_115 | guizhou | 13347420870 | NULL | 23| demo | userinfo | 10007 | user_116 | chengdu | 13347420870 | NULL | 24+---------------+------------+-------+----------+-----------+--------------+-------+ 259 rows in set (0.18 sec)

上述Doris On Iceberg我们只演示了Iceberg单表的查询,你还可以联合Doris的表,或者其他的ODBC外表,Hive外表,ES外表等进行联合查询分析,通过Doris对外提供统一的查询分析入口。

自此我们完整从搭建Hadoop,hive、flink 、Mysql、Doris 及Doris On Iceberg的使用全部介绍完了,Doris朝着数据仓库和数据融合的架构演进,支持湖仓一体的联邦查询,给我们的开发带来更多的便利,更高效的开发,省去了很多数据同步的繁琐工作。

作者:京东零售 吴化斌

来源:京东云开发者社区 转载请注明来源

点赞
收藏

评论区

加载中...

相关推荐

Flink集成数据湖之实时数据写入iceberg

背景iceberg简介flink实时写入准备sqlclient环境创建catalog创建db创建table插入数据查询代码版本总结

Flink 助力美团数仓增量生产

简介:本文由美团研究员、实时计算负责人鞠大升分享,主要介绍Flink助力美团数仓增量生产的应用实践。内容包括:1、数仓增量生产;2、流式数据集成;3、流式数据处理;4、流式OLAP应用;5、未来规划。一、数仓增量生产1.美团数仓架构先介绍一下美团数仓的架构以及增量生产。如下图所示,这是美团数仓的简单架构,我

Flink 1.11 与 Hive 批流一体数仓实践

导读:Flink从1.9.0开始提供与Hive集成的功能,随着几个版本的迭代,在最新的Flink1.11中,与Hive集成的功能进一步深化,并且开始尝试将流计算场景与Hive进行整合。本文主要分享在Flink1.11中对接Hive的新特性,以及如何利用Flink对Hive数仓进行实时化改造,从而实现批流