Beam DataFrame执行数据操作后如何正确转换回PCollection

可正常运行的基础代码
下方代码可直接运行,输出结果符合预期:
table_pcollection = (p | 'Read table' >> io.ReadFromBigQuery(table=f'{TABLE_PREFIX}.test_table')) # 仅保留表中指定列 df_table = to_dataframe(table_pcollection, proxy=pandas.DataFrame(columns=['col1', 'col3', 'col8', 'is_active']), label=f'CNV: To dataframe') test_pcoll = (to_pcollection(df_table, label='to pcollcetion final pcollection', yield_elements='pandas', include_indexes=True)) test_pcoll | 'Output of test_pcoll' >> io.WriteToText('output/test_pcoll.txt')
操作DataFrame后的异常情况
如果对上述转换得到的DataFrame执行where、groupby、drop等操作,代码会运行失败:
table_pcollection = (p | 'Read table' >> io.ReadFromBigQuery(table=f'{TABLE_PREFIX}.test_table')) df_table = to_dataframe(table_pcollection, proxy=pandas.DataFrame(columns=['col1', 'col3', 'col8', 'is_active']), label=f'CNV: To dataframe') df_table.where((df_table['is_active']) & (df_table['col3'] == 'new'), inplace=True) df_table.dropna(inplace=True) test_pcoll = (to_pcollection(df_table, label='to pcollcetion final pcollection', yield_elements='pandas', include_indexes=True)) test_pcoll | 'Output of test_pcoll' >> io.WriteToText('output/test_pcoll.txt')
运行时抛出错误如下:
AttributeError: 'str' object has no attribute 'eq' [while running 'to pcollcetion final pcollection/[ComputedExpression[get_column_Series_4878315232], ComputedExpression[get_column_Series_4878353168], ComputedExpression[eq_Series_4878353600], ComputedExpression[__and___Series_4879336496]]:4879347568/FlatMap(evaluate)/FlatMap(evaluate)']
备注:如果在PCollection转DataFrame前,先增加
beam.Select步骤设置schema,再获取DataFrame对象,就可以正常执行各类DataFrame操作,也能成功转换回PCollection。
问题原因与解决方案
报错核心原因:to_dataframe返回的不是普通本地Pandas DataFrame,而是Beam实现的延迟计算(Deferred)DataFrame,它只会构建计算表达式树,不会立刻执行操作,因此完全不支持inplace=True这类原地修改参数——原地修改会直接破坏表达式树结构,导致后续to_pcollection执行计算时,把表达式节点误识别为普通字符串对象,最终抛出属性错误。
修复方式非常简单,去掉所有inplace=True参数,将操作返回的新DataFrame重新赋值给变量即可:
table_pcollection = (p | 'Read table' >> io.ReadFromBigQuery(table=f'{TABLE_PREFIX}.test_table')) df_table = to_dataframe(table_pcollection, proxy=pandas.DataFrame(columns=['col1', 'col3', 'col8', 'is_active']), label=f'CNV: To dataframe') # 移除inplace=True,赋值操作返回的新对象 df_table = df_table.where((df_table['is_active']) & (df_table['col3'] == 'new')) df_table = df_table.dropna() test_pcoll = (to_pcollection(df_table, label='to pcollcetion final pcollection', yield_elements='pandas', include_indexes=True)) test_pcoll | 'Output of test_pcoll' >> io.WriteToText('output/test_pcoll.txt')
如果仍出现类型相关报错,可以配合备注提到的beam.Select提前明确PCollection的Schema,避免Beam自动推断字段类型出错。
内容的提问来源于stack exchange,提问作者Sanket
相关产品推荐
相关产品推荐

