如何将S3中TXT文件的SQL语句转为PySpark可执行查询语句
解决方案
你可以通过以下步骤将读取到的Spark DataFrame转换为完整的SQL字符串,再传入Redshift查询代码中:
完整实现代码
# 读取S3上的SQL文件 hist_sql = spark.read.text('s3://team-test/history/sql/his.txt') # 提取所有行文本并拼接成完整SQL字符串 sql_lines = hist_sql.select("value").rdd.flatMap(lambda x: x).collect() sql = "\n".join(sql_lines) # 可选:清理原SQL末尾多余的句号(匹配你的示例需求) sql = sql.rstrip('.').strip() # 执行Redshift查询 df = spark.read.format("com.databricks.spark.redshift") \ .option("url", jdbcUrl) \ .option("query", sql) \ .option("forward_spark_s3_credentials", True) \ .load()
代码说明
spark.read.text生成的DataFrame默认只有一列value,每行对应原SQL文件的一行内容select("value").rdd.flatMap(lambda x: x).collect()将DataFrame中的每行文本提取为字符串列表"\n".join(sql_lines)用换行符拼接列表元素,还原原SQL的换行格式rstrip('.').strip()用于去除原SQL末尾多余的句号(如示例中的condition2.),同时清理首尾空白字符,可根据实际文件内容调整或省略
内容的提问来源于stack exchange,提问作者MLDL
相关产品推荐
相关产品推荐

