纵腾湖仓全链路落地实践
以下文章来源于 DataFunTalk,作者陈晶

DataFunTalk
专注于大数据、人工智能技术应用的分享与交流。致力于成就百万数据科学家。定期组织技术分享直播,并整理大数据、推荐/搜索算法、广告算法、NLP自然语言处理算法、智能风控、自动驾驶、机器学习/深度学习等技术应用文章。
导读
数据湖作为一个统一存储池,可接入多种方式的数据输入,无缝对接多种计算分析引擎,进行高效的数据处理与分析。本文将介绍数据湖上的选型思考与探索实践。
主要内容包括以下四个部分:
- 总体架构
- 入湖方案选型
- 实时入湖优化
- 数据湖上的查询
01 总体架构
面对日益增长的数据量,Lambda架构使用离线/实时两条链路和两种存储完成数据的保存和处理。这种繁杂的架构体系带来了不一致的问题,需要通过修数、补数等一系列监控运维手段去弥补。为了统一简化架构,提高开发效率,减少运维负担,我们实施了基于数据湖Hudi+Flink的流批一体架构,达到了降本增效的目的。
如下图所示,总体架构包括数据采集、ETL、查询、调度、监控、数据服务等。要解决的是数据从哪里来到哪里去,怎么过去,怎么用,以及过程中的调度和监控、元数据管理、权限管理等问题。

- 数据从哪里来:我们的数据来自MySQL、MongoDB、Tablestore、Hana。
- 数据到哪里去:我们的数据会写入到Hudi、Doris,其中Doris负责存储部分应用层的数据。
- 数据怎么过去:将在后面的实时入湖部分进行介绍。
- 数据用在哪里:我们的数据会被OLAP、机器学习、API、BI查询使用,其中OLAP和BI都通过Kyuubi的服务进行查询。
- 任务调度:主要通过DolpuinScheduler来执行,基于quartz的cronTrigger完成shell、SQL等调度。
- 监控部分:通过Prometheus和Grafana,这是业界通用的解决方案。
- 元数据采集:通过DataHub完成,采用了datahub的ingestion framework框架来采集各种数据源的元数据。
- 权限管理:主要包括Kyuubi服务端的统一认证和引擎端的独立鉴权。
02 入湖方案选型
数据入湖方案设计上,我们比较了三种入湖的实现思路。
1. 入湖方案一
如下图所示,包含了两条支线:
- 分支①:Flink SQL通过MySQL-CDC connector和Hudi connector完成source和sink端读写。这样MySQL每张表由单独的binlog dump线程读取binlog。
- 分支②:通过MySQL多库表配置一个Debezium Connector实现单独binlog dump线程读取多库表,解析后发送到Kafka的多个topic。即一张表一个topic。之后用Flink SQL通过Kafka connector和Hudi connector完成source和sink端读写。
这种方案的主要优点是Flink和CDC组件都经过了充分验证,已经非常稳定成熟了。而主要缺点是Flink SQL需要定义表DDL。但我们已经开发DDL列信息从元数据系统获取,无须自定义。并且写Hudi是每张表一个Flink任务,这样会导致资源占用过多。另外Flink CDC还不支持Schema演变,一旦Schema变更,需要重新拉取数据。

2. 入湖方案二
这一方案是在前一个方案分支二的基础上进行了一定的改进,通过Dinky完成整库数据同步,其优点是同源数据合并成一个source节点,减轻源库压力,根据schema、database、table分流sink到对应表。其缺点是不支持schema演变,表结构变更须重新导数。如下图所示,mysql_biz库中有3张表,从flink dag图看到mysql cdc source分3条流sink到Hudi的3张表。

3. 入湖方案三
主要流程如下图所示。其主要优点是支持Schema演变。Schema变更的信息由Debezium注册到Confluence Schema Registry,schema change的信息通过DeltaStreamer执行任务变更到Hudi,使得任务执行过程中不需要重新拉起。其主要缺点是依赖于Spark计算引擎,而我们部门主要用Flink,当然,这会因各个公司实际情况而不同。
下图分别是Yarn的deltastreamer任务,Kafka schema-change topic的DML message和Hudi表变更后的数据。

4. 入湖方案总结
在方案选型时,可以根据下面的流程图进行比较选择:

- 先看计算框架是Spark还是Flink,如果是Spark则选择方案三,即Deltastreamer,这一方案适用于表结构变更频繁,重新拉取代价高,主要技术栈是Spark的情况。
- 如果是Flink,再看数据量是否较少,以及表结构是否较稳定,如果是的话,选择方案二,Dinky整库同步方案支持表名过滤,适用数据量较少且表结构较稳定的表。
- 如果否,再考虑mysql能否抗较大压力,如果否,那么选择方案一下分支,即Kafka Connect,Debezium拉取发送Kafka,从Kafka读取后写Hudi。适用数据量较大的多张表。
- 如果是,则选择方案一上分支,即Flink SQL mysql-cdc写Hudi,适用于对实时稳定要求高于资源敏感的重要业务场景。
03 实时入湖优化
我们的入湖场景是Flink Stream API读取Pulsar写Hudi MOR表,特点是数据量大,并且源端的每条消息都只包含了部分的列数据。我们通过使用Hudi的MOR表格式和PartialUpdateAvroPayload实现了这个需求。使用Hudi的MOR格式,是因为COW的写放大问题,不适合数据量大的实时场景,而MOR是增量数据写行存Avro格式log,通过在线或离线方式压缩合并至列存格式parquet。在保证写效率的同时也兼顾了查询的性能。不过需要通过合并任务定期地对数据进行合并处理,这是引入复杂度的地方。
以下面这张图为例,recordKey是ID1的3条msg,每条分别包含一个列值,其余字段为空,按ts列precombine,当ts3 > ts2 > ts1时,最终Hudi存的ID1行的值是v1,v2,v3,ts3。

此入湖场景痛点包括,MOR表索引选择不当,压缩异常导致越写越慢,直至checkpoint超时,某分区存在重复文件导致写任务出错,MOR表某个压缩计划pending阻碍此bucket的压缩及后续的压缩计划生成,以及如何平衡效率与资源等。我们在实践过程中针对一些痛点实施了相应的解决方案。
- Hudi表索引类型选择不当,导致越写越慢至CK超时,这是因为Bucket索引通过hash映射recordKey到fileGroup。而Bloom索引是保存recordKey和partition、fileGroup值来实现,因此checkpoint size会随数据量的增加而增长。Bloom Filter索引基于布隆过滤器实现,索引信息存储在parquet的footer中,Bloom的假阳性问题也会导致更新越来越慢,假阳性是指只能判断数据一定不在某个文件而不能保证数据一定在某个文件,因此存在多个文件都可能存在某条数据,即须读取多个文件才能准确判断。我们做的优化是使用Bucket索引代替Bloom索引,Hudi目前也支持了可以动态扩容的Bucket参数。

- MOR表压缩执行异常,具体来说有以下三个场景:
- 单log超过1G,使写延迟提高,导致越写越慢至checkpoint超时,checkpoint端到端耗时增长至3-6分钟;
- 在inline schedule的压缩模式下,offline execute出现报错:log文件不存在;
- Compaction一直处于Infight状态,即进行中,不能完成;同时存在无效compaction,既不能被压缩,也不能被取消。
此3种现象的原因都是Sink:compact_commit算子的并行度 > 1,我们做的优化是降低压缩过程的并发度,设置compact_commit Parallelism = 1。并行度改成1后1G的log压缩正常。整张表size明显减少。log到parquet的压缩比默认是0.35。

- MOR表某分区存在重复文件,导致写任务出错。出现这个问题的原因是某个instant已写log文件但未成功提交到timeline时,发生异常重启后未rollback这个instant,即未清理已有log,继续写此instant则有重复。我们做的优化是在遇到重复文件时,通过Hudi-Cli执行去重任务,再恢复执行。具体来说,需要拆分成以下四个步骤:
- 停止当前的Flink任务;
- 通过Hudi-cli执行去重命令;
- 删除partition文件,修复文件移到原分区;
- 重新启动Flink任务。

- MOR表某个压缩计划pending,阻碍此bucket的压缩及后续的压缩计划生成。这个问题是由于环境问题导致的zombie compaction或bug。上图中第一列是compaction instant time,即压缩计划生成时间,第二列是状态,第三列是此压缩计划包含的文件数。8181的instant卡住,且此压缩计划包含2198个文件,即涉及到大量的file group,涉及的file group不会有新的压缩计划生成。导致表的size增加,写延时。我们做的优化是回滚不正常的合并任务,重新处理。即利用较多资源快速离线压缩完。保证之后启动的Flink任务在相对少的资源情况下仍然可以保证更新和在线压缩的效率。具体来说,包括下面的命令:
- 执行HoodieFlinkCompactor把所有inflight instant回滚成requested状态;
- 执行compaction unschedule命令。

经过多次的修改和验证,我们的入湖任务在性能和稳定性上取得了明显的改善。在稳定性上,做到了在十几天内任务无异常。在时延上,做到了分钟级别的checkpoint和数据可见。在资源使用上,对Hadoop YARN资源的占用明显减少。

下图总结了我们对实时入湖做的参数优化方案,包括:
- 索引选择GLOBAL_BLOOM -> BUCKET_INDEX #Bucket索引较布隆索引写吞吐性能高;
- BUCKET_NUM 20 #Bucket数量:根据单分区数据量评估,保证File Slice2GB,平衡读写性能;
- Flink增量checkpoint:Rockdb #Flink ck存储,rockdb支持增量ck,减少单ck数据量,提高写吞吐;
- Yarn资源:jobmanager 5G #Flink jobmanager内存,减少oom,保证稳定;taskmanager 50G 20S #Flink taskmanager内存与slot数,slot与并发度、bucket数一致;
- write.rate.limit 30000 #写速度限制,过载保护,保证作业稳定运行;
- write.max.size 2560 #写用到最大内存,于taskmanager每个slot内存一致;
- write.batch.size 512 #批量写,适量调大减少刷盘频率;
- compaction.max.memory 2048 #压缩用到的最大内存,适量调大提升压缩速度;
- compaction.trigger.strategy num_and_time #压缩策略 增量提交个数或时间达标触发生成压缩策略;
- compaction.delta_seconds 30 #压缩策略之时间,减少时间间隔,减少单个压缩文件数;
- compaction.delta_commits 2 #压缩策略之增量提交个数,减少个数,减少单个压缩文件数。

实时任务入湖的优化思路流程包括下面几个步骤:
- 先确定bucket数量,观察fileSlice大小,估算调整;
- 根据bucket数确定Flink job并行度,与bucket数保持一致;
- 根据并行度确定tm资源,即并行度 = 总slot数;
- 根据总slot数确定内存,即总内存 = num(slot) * (write.max.size);
- 根据pulsar topic流量确定write.rate.limit,一般峰值 * 1.5;
- 根据Flink job内存使用情况及平稳度确定write.max.size,可拿追存量数据测试,一般内存降到写速度明显下降即为内存最低值;
- 根据write.max.size确定write.batch.size和compaction.max.memory,前者是后两个的和;
- 根据pulsar topic流量确定压缩类型,一般超过10w/s考虑使用inline schedule和offline execution。
04 数据湖上的查询
在引入Kyuubi前,我们通过JDBC、Beeline、Spark Client、Flink Client等客户端访问服务层执行查询,没有统一入口,多个平台不互通,多账号权限体系。用户的痛点是跨多平台开发体验差,低效率。平台层的痛点是问题定位运维复杂,存在资源浪费。
在引入Kyuubi后,我们基于社区版Kyuubi做了一定的改造,包括JDBC引擎开发、JDBC引擎Ranger鉴权开发、BI、JDBC客户端元数据适配修改、Spark引擎大结果集存HDFS、支持导数开发、JDBC引擎SQL拦截控流开发等,实现了统一数据服务入口,做到了统一认证权限管理和统一易用原则。
下图展示了Kyuubi的架构和权限管控:

Kyuubi查询流程是:客户端请求通过LDAP认证后,连接Kyuubi Server生成Kyuubi session,之后Kyuubi server根据连接用户以及用户隔离级别路由到已经启动的engine或启动一个新的engine。Spark引擎会先申请container运行AppMaster,后申请container运行executor执行task。Flink引擎会完成StreamGraph至JobGraph至executionGraph构建并通过Jobmanager和taskmanager运行。其中engine端RangerPlungin会在SQL解析后拉取RangerAdmin由用户配置的策略进行鉴权。RangerAdmin完成用户同步,策略刷新等。
Kyuubi on Flink跨库查询的目的是尝试基于Flink实现流批一体,支持跨数据源导数SQL化。我们的实现方案是通过Flink Metadata Catalog Connector的开发,即基于元数据系统以统一datasource.db.table的格式查询所有数据源,且让用户免于自定义DDL。其中元数据采集是用datahub的ingestion framework采集各种数据源的元数据,并生成对应Flink表属性。Flink端是扩展AbstractCatalog查询metadata DB,实现CatalogFactory接口。
其基本流程如下图所示:

完整流程是:
- 发起采集请求;
- 和3是采集服务调Datahub ingestion framework完成元数据采集并写到metadata DB同时写Flink表属性;
- 用户发送SQL到Kyuubi server;
- Kyuubi server发送SQL到Flink engine;
- 和6是Flink metadata catalog会读取metadata DB根据Flink表属性读取对应数据源。
Kyuubi on JDBC Doris可以通过外表查询Hudi,但在Doris 1.2版本,仍然有一定的限制,Hudi目前仅支持Copy On Write表的Snapshot Query,以及Merge On Read表的Read Optimized Query。后续将支持Incremental Query和Merge On Read表的Snapshot Query。
Doris的架构示意和其基本使用流程如下图所示:

今天的分享就到这里,谢谢大家。

end
