Apache Beam SqlTransform单节点运行引发Dataflow拖慢问题排查与解决
在Dataflow管道中使用SqlTransform组件计算滚动平均值等窗口聚合操作,管道配置最多可使用3个Worker,但执行时该组件仅用1个Worker,还出现“拖慢任务(straggler)”,处理速度极慢。经排查问题出在NUM_1 / SUM(NUM_1) OVER() AS COST_2_DATASET_AVERAGE这行计算语句。
咨询问题
- 该步骤为何无法使用多个Worker?
- 如何解决此问题?
附使用的窗口查询语句:
windowing_query = """SELECT DATE_STR, SUBS_ID, (NUM_1 + NUM_2 + NUM_3) AS TOTAL_COST, NUM_1 / SUM(NUM_1) OVER() AS COST_2_DATASET_AVERAGE, AVG(NUM_1) OVER (w ROWS 2 PRECEDING) as AVERAGE_COST_LAST_3MTHS FROM PCOLLECTION WINDOW w AS (PARTITION BY SUBS_ID ORDER BY DATE_STR)""" feature_rows = rows_train_dataset | SqlTransform(windowing_query)
1. 无法使用多Worker的原因
SUM(NUM_1) OVER()是无分区、无框架的全局窗口聚合,它需要计算整个数据集所有NUM_1的总和。这种全局聚合逻辑上要求所有数据都集中到同一个节点(Worker)进行计算——因为只有拿到全量数据才能算出准确的全局总和,Dataflow无法将这个计算任务拆分到多个Worker并行执行,所以只能用1个Worker处理,自然会出现单节点负载过高、速度慢的拖慢任务问题。
而你同时使用的AVG(NUM_1) OVER (w ROWS 2 PRECEDING)是按SUBS_ID分区的滚动窗口,这个是可以并行的,但全局聚合的存在会导致整个SqlTransform步骤被强制串行执行。
2. 解决方法
方法一:预计算全局总和并广播(推荐,适合精确计算场景)
把全局SUM(NUM_1)的计算从SqlTransform中拆分出来,单独做全局聚合后广播给所有Worker,再和原数据集关联后计算比值:
- 先计算全局NUM_1的总和:
from apache_beam.transforms.combiners import CombineGlobally # 计算全局NUM_1总和 global_sum_num1 = rows_train_dataset | CombineGlobally(lambda elements: sum(e['NUM_1'] for e in elements)).As('global_sum_num1') - 将全局总和广播,和原数据集合并:
from apache_beam.transforms.util import Broadcast, CoGroupByKey # 给原数据集加一个固定key,方便和广播的全局值关联 keyed_rows = rows_train_dataset | 'Add fixed key' >> Map(lambda x: ('fixed_key', x)) # 广播全局总和并加相同key broadcast_sum = global_sum_num1 | 'Broadcast sum' >> Broadcast() | Map(lambda x: ('fixed_key', x)) # 合并数据集 merged = (keyed_rows, broadcast_sum) | CoGroupByKey() # 展开合并后的数据,把全局总和和每条数据关联 joined_rows = merged | 'Expand merged data' >> Map( lambda x: {**x[1][0][0], 'GLOBAL_SUM_NUM1': x[1][1][0]} ) - 修改SqlTransform的查询语句,用预计算的全局总和替换窗口函数:
windowing_query = """SELECT DATE_STR, SUBS_ID, (NUM_1 + NUM_2 + NUM_3) AS TOTAL_COST, NUM_1 / GLOBAL_SUM_NUM1 AS COST_2_DATASET_AVERAGE, AVG(NUM_1) OVER (w ROWS 2 PRECEDING) as AVERAGE_COST_LAST_3MTHS FROM PCOLLECTION WINDOW w AS (PARTITION BY SUBS_ID ORDER BY DATE_STR)""" feature_rows = joined_rows | SqlTransform(windowing_query)
这样全局总和的计算可以并行(CombineGlobally会自动拆分计算局部和再汇总),后续的SqlTransform步骤也能按SUBS_ID分区并行执行,充分利用多个Worker。
方法二:使用近似聚合(仅适用于允许误差的场景)
如果业务场景允许结果存在一定误差,可以使用Dataflow支持的近似求和函数(如APPROX_SUM),这类函数可以并行计算,不需要集中全量数据:
windowing_query = """SELECT DATE_STR, SUBS_ID, (NUM_1 + NUM_2 + NUM_3) AS TOTAL_COST, NUM_1 / APPROX_SUM(NUM_1) OVER() AS COST_2_DATASET_AVERAGE, AVG(NUM_1) OVER (w ROWS 2 PRECEDING) as AVERAGE_COST_LAST_3MTHS FROM PCOLLECTION WINDOW w AS (PARTITION BY SUBS_ID ORDER BY DATE_STR)"""
注意:近似聚合的结果不是精确值,需要根据业务需求判断是否适用。
方法三:拆分全局聚合为两步并行计算
先在每个分区计算局部SUM,再汇总全局SUM,最后关联回原数据:
- 先做局部聚合(按任意可并行的键分区,比如SUBS_ID或者日期分区):
partial_sums = rows_train_dataset | 'Partial sum by SUBS_ID' >> CombinePerKey(lambda elements: sum(e['NUM_1'] for e in elements)) - 再计算全局总和:
global_sum = partial_sums | 'Global sum' >> CombineGlobally(lambda elements: sum(e[1] for e in elements)).As('global_sum_num1') - 后续步骤和方法一类似,广播全局总和并关联原数据,再执行SqlTransform。这种方式和方法一本质类似,但拆分了局部和全局的聚合步骤,同样能利用多Worker并行。
内容的提问来源于stack exchange,提问作者cr

