如何在Apache Beam中引用SQL文件而非直接使用字符串SQL?
解决方案
问题根源
你读取SQL失败的核心原因是**.sql文件里的SQL语句多了一层外层双引号**,导致读取后的sql_query字符串包含多余引号,BigQuery无法识别这种格式的SQL语句。
步骤1:修正SQL文件内容
把file_query.sql里的内容改成纯SQL语句,去掉外层双引号:
SELECT col1, col2, col3 FROM `project.dataset.table`
步骤2:正确读取SQL文件
保留你原来的读取逻辑,建议加上strip()移除文件末尾的换行或空格,同时确保文件路径正确(如果不在当前工作目录,需用绝对路径):
with open("file_query.sql", "r") as myfile: sql_query = myfile.read().strip()
步骤3:修正原代码的变量错误
你代码里定义了apromore_event_schema,但写入Parquet时用的是未定义的event_schema,需要统一变量名:
process_data = (pipeline | "ReadFromBigQuery" >> beam.io.Read(bq_source) | "Write to Parquet" >> beam.io.parquetio.WriteToParquet( file_path_prefix='gs://project-gcs/test/', schema=apromore_event_schema, # 替换为正确的变量名 file_name_suffix='.parquet' ) )
完整修正后的代码片段
class SimplePipeline: def run_pipeline(self, known_args, pipeline_args, streaming_flag): pipeline_options = PipelineOptions(pipeline_args, streaming=streaming_flag, save_main_session=True) custom_options = pipeline_options.view_as(self.get_pipeline_options_class()) log_info("Starting Pipeline") # 读取外部SQL文件 with open("file_query.sql", "r") as myfile: sql_query = myfile.read().strip() with Pipeline(options=pipeline_options) as pipeline: bq_source = beam.io.BigQuerySource(query=sql_query, use_standard_sql=True) apromore_event_schema = pa.schema([ pa.field("col1", "string", metadata={'parent':'id'}), pa.field("col2", "string", metadata={'parent':'name'}), pa.field("col3", "string", metadata={'parent':'age'})] ) process_data = (pipeline | "ReadFromBigQuery" >> beam.io.Read(bq_source) | "Write to Parquet" >> beam.io.parquetio.WriteToParquet( file_path_prefix='gs://project-gcs/test/', schema=apromore_event_schema, file_name_suffix='.parquet' ) )
内容的提问来源于stack exchange,提问作者dragon
相关产品推荐
相关产品推荐

