【发布时间】:2022-12-20 11:45:07
【问题描述】:
我是大数据的新手。我有几个大表(~TB 规模),其中包含不同月份的数据,我正在尝试分析这些表的特征漂移。我专门尝试连续两个月计算相应列的 KL 散度。在计算 KL 散度之前,我需要获得列的概率分布,这意味着创建一个直方图,我可以在其中包含 bin 和计数。归一化计数数组将为我提供使用 scipy 熵函数计算 KL 散度所需的概率分布。
我正在分析的表有数百万行和大约 2000 列/特征,并且都在 BigQuery 中。我尝试使用两种不同的工具来解决这个问题。
(我的所有方法都使用 Python)
1- 我尝试使用 Pyspark 并且只需要 70 秒来计算一张表的一列的 bins 和计数。以这种方式,我需要数周时间才能完成我拥有的所有功能和表格。
2- 我利用大查询 python api 并创建了 python 函数来批量创建长查询(例如 10 列的批次)来计算每列的 bins 和计数。为了使用大查询计算 bin 和计数,我使用了 bigquery 的“CASE WHEN”功能并将我的值设置为预定义的 bin 中心。下面是一个例子
case when col_name1>=1 and col_name1<2 then bin_center_array[0]
when col_name1>=2 and col_name1<3 then bin_center_array[1]
...
使用大查询,计算每列只需要 0.5 秒(整个计算不到 2 小时,而不是一周)。但是,如果我在两个表上执行 10 个批次,我将在大约 10 个批次后用完 QueryQuotaPerDayPerUser(请注意,我需要 2000/10=200 个批次)。如果我将批处理大小增加到更大的值,我会得到“BadRequest:超过 400 个资源......”错误(注意:每个批处理本质上都会产生一个长查询,批处理越大,查询越长)。
我不确定如何解决这个问题。任何帮助或建议表示赞赏
【问题讨论】:
-
一种可能的快速绕道是采用采样方法,例如 FARM_FINGERPRINT 或 TABLESAMPLE SYSTEM。
-
增加并发批处理查询的 quota Limit 是否有助于您的设置?
标签: python pyspark google-bigquery bigdata