一、大规模数据场景下元数据采集的核心痛点
很多公司在搭建数据中台的时候,都会用到DataHub来管理各种元数据,但实际用到大规模数据系统(比如几百个数据源、几十亿条元数据)的时候,就会发现采集速度变得特别慢,小数据量的时候跑采集任务可能几分钟就完了,到了大规模场景下,不仅要跑十几个小时,还经常因为超时或者服务压力大中断,这就是典型的元数据采集性能瓶颈。这个问题其实不是DataHub本身不行,而是很多开发者在使用的时候没注意到大规模场景下的特殊需求,导致原本好用的工具变“卡”了。
1.1 小数据与大规模场景的差异
如果你的公司只有几十个数据源,每个数据源的元数据也就几万条,那随便用DataHub默认的配置就足够了,根本不会有性能问题。但如果你的公司有上百个数据源,每个数据源的表就有几千张,每张表还有字段、权限、血缘关系等元数据,那时候再用默认的串行采集方式,就会发现每采集一个数据源就要等很久,加起来的时间完全无法接受。
1.2 常见的痛点表现
除了采集时间太长之外,还有几个明显的问题:比如采集过程中DataHub服务经常报超时错误,或者突然内存占用飙升到100%,甚至连服务都崩了;还有采集到一半中断,需要重新跑,浪费时间;另外,即使采集完了,查询元数据的时候也会变慢,这其实和采集阶段的性能瓶颈是相关的。
二、性能瓶颈的深度拆解
要解决问题,得先搞清楚瓶颈到底在哪里,很多时候都是几个点共同作用导致的,不是单一原因。
2.1 采集端的串行执行瓶颈
最常见的就是所有数据源的采集任务是一个个按顺序跑的,就像你要去10个不同的超市买东西,一个超市一个超市逛,总时间肯定长。DataHub默认的采集方式很多是串行的,每个采集任务要等前一个完了才会启动,大规模情况下,光等待的时间就占了大部分。
2.2 元数据存储的单点压力
DataHub的元数据通常存在像Elasticsearch或者MongoDB这样的存储里,默认配置下,写入操作是单点的,也就是所有采集来的元数据都往同一个节点写,当每秒写入几万条的时候,这个节点的读写压力就会爆炸,导致写入超时,采集中断。
2.3 批量处理的阈值问题
很多人不知道DataHub的采集有个批量大小的配置,默认可能是每次传100条元数据过去,也就是每次和服务端通信只发100条。如果改成每次传1000条,那通信的次数就会减少很多,减少网络开销,速度自然会变快,但很多人没调整这个参数,导致每次通信都要花时间,慢就不奇怪了。
三、针对性突破方案与示例
针对上面的瓶颈,我们可以从三个方面入手:把串行改成并行,优化存储的写入方式,调整批量处理的阈值。下面用实际的Python代码示例来展示,技术栈用Python 3.9 + DataHub SDK 0.11.0,都是很容易上手的工具,不用复杂的环境。
3.1 示例前置说明
这个示例的目标是对比串行采集和并行采集的时间差异,并行采集就是把多个数据源的任务同时跑,这样总时间就和最慢的那个任务差不多,而不是累加。我们用的DataHub REST Emitter是官方提供的,用来把元数据发送到DataHub服务端。
3.2 原始的串行采集代码
from datahub.emitter.rest_emitter import DatahubRestEmitter
from datahub.metadata.schema_classes import DatasetUsageStatsClass
import time
# 初始化DataHub连接,替换成你自己的服务地址
emitter = DatahubRestEmitter(gms_server="http://你的DataHub地址:8080")
# 模拟10个数据源,实际场景这里换成真实的数据源列表
data_sources = [f"mysql_source_{i}" for i in range(10)]
def collect_single_source(source_name):
"""单个数据源的采集函数,这里模拟耗时操作,实际是采集元数据的逻辑"""
time.sleep(0.5) # 模拟每个数据源采集需要0.5秒
# 构造元数据,实际场景这里会从数据源读取真实的元数据
usage_stats = DatasetUsageStatsClass(
dataset=f"urn:li:dataset:(urn:li:dataPlatform:mysql,{source_name},PROD)",
totalQueries=10000 + i*500 # 模拟查询次数
)
emitter.emit(usage_stats)
print(f"数据源{source_name}采集完成")
# 串行执行所有采集任务
start_time = time.time()
for source in data_sources:
collect_single_source(source)
serial_total_time = time.time() - start_time
print(f"串行采集总时间:{serial_total_time:.2f}秒")
这段代码的执行时间大概是5秒,因为10个任务每个0.5秒,加起来刚好5秒,在大规模场景下,几百个数据源就是几百秒,时间完全不可控。
3.3 优化后的并行采集代码
from concurrent.futures import ThreadPoolExecutor, as_completed
import time
# 最大线程数,根据你的DataHub服务负载调整,这里设为5就不会太占资源
MAX_THREADS = 5
start_time = time.time()
# 用线程池执行并行任务
with ThreadPoolExecutor(max_workers=MAX_THREADS) as executor:
# 把所有采集任务提交给线程池
futures = [executor.submit(collect_single_source, source) for source in data_sources]
# 等待所有任务完成,同时可以做进度跟踪或异常处理
for future in as_completed(futures):
pass
parallel_total_time = time.time() - start_time
print(f"并行采集总时间:{parallel_total_time:.2f}秒")
print(f"时间节省了:{serial_total_time - parallel_total_time:.2f}秒")
这段代码的执行时间大概是1.2秒左右,因为同时跑5个任务,0.5秒完成5个,剩下5个再跑0.5秒,总共1秒多,比串行快了4倍左右,实际场景下,几百个数据源的话,时间节省会更明显。
3.4 其他优化小技巧
除了并行采集,还有两个小技巧可以提升性能:一个是调整DataHub的批量大小参数,比如把batch_size从默认的100改成500或者1000,减少网络请求次数;另一个是用批量写入的API,DataHub有专门的批量emit方法,比单条发送更快,比如emitter.emit_many(),一次传多条元数据,减少网络开销。
四、应用场景、技术优缺点与注意事项
4.1 核心应用场景
这种性能突破方案非常适合这些场景:比如电商平台的全链路元数据采集,涉及用户订单、商品、支付等上百个数据源;金融行业的交易数据、用户画像数据的元数据管理,数据源多且元数据量大;大数据平台的Hive、Spark、HBase等组件的元数据统一采集,需要处理几百张以上的表;还有政府或者大型企业的核心数据中台,需要管理全公司所有业务系统的元数据,规模大,采集频率高。
4.2 技术优缺点
优点很明显:一是采集时间大幅缩短,从小时级降到分钟级,适合大规模场景;二是简单易实现,不用改DataHub的核心代码,只要用SDK加几行并行的代码就可以;三是兼容性好,不管你用的是DataHub社区版还是企业版,都能生效。缺点的话,就是并行任务的数量需要调整,不能设太多,不然会超过DataHub服务的负载,导致服务崩溃;还有需要监控采集的过程,避免出现单个任务失败没被发现,导致元数据不全。
4.3 注意事项
第一个注意事项是线程数的设置,一定要根据你的DataHub服务的配置来,比如你的DataHub服务是4核8G的,那线程数设为3或者4就够了,设太多会导致服务压力太大,采集反而变慢;第二个是批量大小的设置,不能太大,比如设到10000的话,单个请求的 payload 太大,网络传输慢,也会出问题,建议先从500开始试,再根据实际情况调整;第三个是要做采集任务的重试,比如某个数据源采集失败,要自动重试几次,避免因为临时网络问题或者数据源服务波动导致整个采集中断;第四个是要定期清理不需要的元数据,比如已经下线的数据源的元数据,不然存储压力会越来越大,间接影响采集性能。
五、总结
DataHub在大规模数据系统中的元数据采集瓶颈,本质上是没有适配大规模场景下的并发需求和负载特性,通过把串行采集改成并行,调整批量处理的参数,优化写入的方式,就可以大幅提升采集的性能。这个方案门槛低,只要懂一点Python的基础就能实现,适合大部分用DataHub管理元数据的团队。当然,如果你遇到特别大规模的场景,比如上万个数据源,那可能还要用分布式采集的架构,但对于绝大多数大型企业的日常元数据管理需求来说,上面的优化方案已经足够解决问题,让DataHub能从容应对大规模数据系统的元数据采集任务。
评论
围绕“DataHub在大规模数据系统中元数据采集的性能瓶颈与突破”参与讨论