Beam DirectRunner Calcite SqlTransform无法指定自定义表名报错问题
问题原因
你遇到的Object 'my_table' not found报错是因为Apache Beam Python SDK的SqlTransform不支持直接将{自定义表名: PCollection}的字典通过管道符|传入的方式注册输入表,该语法仅在Java SDK中生效。
单输入场景下直接传入PCollection对象时,Beam会默认将该输入的表名设为PCOLLECTION,所以你将SQL中的表名改为PCOLLECTION后可以正常运行,这是单输入的简化适配逻辑。
解决方案
如果需要使用自定义表名,需要调用SqlTransform提供的from_方法显式注册输入表,示例代码如下:
import apache_beam as beam from apache_beam.transforms.sql import SqlTransform with beam.Pipeline() as p: rows = (p | beam.Create([ beam.Row(col1="val1", col2="col2_val1"), beam.Row(col1="val2", col2="col2_val2"), ])) # 自定义表名写法 (SqlTransform.from_("my_table", rows) | SqlTransform("""SELECT * FROM my_table""") | beam.Map(print))
如果需要关联多个表,可链式调用from_方法注册多张表,示例如下:
(SqlTransform .from_("table_a", pcollection_a) .from_("table_b", pcollection_b) | SqlTransform("""SELECT a.*, b.* FROM table_a a JOIN table_b b ON a.id = b.a_id"""))
内容的提问来源于stack exchange,提问作者steeling
相关产品推荐
相关产品推荐

