【发布时间】:2020-07-12 07:26:30
【问题描述】:
我在 pyspark 中编写了代码,任务是查找与前一天 parquet 文件相比的 Delta/更改记录,并将其写入 csv 并带有标记新记录或现有记录的附加列。这是通过连接列并应用 base64(encode(columns here)) 将其命名为 hash_key 来完成的。
-
读取 9 个 parquet 文件并针对它们注册临时表。
- 对其应用查询,将输出命名为 df2。
- 从 S3 读取前一天的 parquet 文件并将其命名为 df1 并获取它的计数。
- 针对两个数据帧将临时表注册为 t1 和 t2。
-
现在我使用 sqlContext.sql 查找 t2 与 t1 相比的更改记录,使用 hash_key 即 df3 并获取行数。
-
将不同的保存在 csv 文件中,并将 parquet 文件替换为 df2 中的数据。
到目前为止,一切都运行良好。但后来的任务是显示增量记录中哪些行是全新的,哪些是旧行,其中一些列已更改。在这里,我尝试使用 PySpark 代码并加入了两个数据帧 df1 和 df3。但我收到以下错误 GC 开销限制超出 或 超时
作为替代方案,我尝试针对数据帧注册临时表并对其执行 sql 查询。 但问题仍然存在。我真的被它困住了。不知道为什么会消耗资源。
配置:
conf = (SparkConf()
.setAppName("GD_Regex")
.set("spark.executor.instances", "3")
.set("spark.executor.cores", "3")
.set("spark.sql.parquet.enableVectorizedReader", "false")
.set("spark.executor.memory", "4g")
.set("fs.s3a.server-side-encryption-algorithm", "AES256")
.set("spark.sql.autoBroadcastJoinThreshold","-1")
)
使用 hash_key 发现差异
s3_parquet_file = sqlContext.read.parquet(parquet_path)
historical_file_row_count = s3_parquet_file.count()
s3_parquet_file.registerTempTable("s3_parquet_file_temp")
query_response.registerTempTable("query_response_temp")
filtered_data = sqlContext.sql(
"""select * from query_response_temp where hash_key NOT IN ( SELECT hash_key FROM s3_parquet_file_temp )""")
方法一:
ff = filtered_data.alias('df2').join(s3_parquet_file.alias('df1'), ( filtered_data.CTMS_STUDY_NUMBER == s3_parquet_file.CTMS_STUDY_NUMBER ) & ( filtered_data.CTMS_SITE_NUMBER == s3_parquet_file.CTMS_SITE_NUMBER )
& (filtered_data.FULL_NAME_OF_PI == s3_parquet_file.FULL_NAME_OF_PI ) & ( filtered_data.PRIMARY_ROLE_NAME == s3_parquet_file.PRIMARY_ROLE_NAME) & ( filtered_data.CENTER_NAME == s3_parquet_file.CENTER_NAME ) , "outer" )\
.select('df2.*', 'df1.CTMS_STUDY_NUMBER')\
.withColumn('status', when(s3_parquet_file.CTMS_STUDY_NUMBER.isNotNull() & filtered_data.CTMS_STUDY_NUMBER.isNotNull(), 'existing').otherwise('new'))\
.filter(filtered_data.CTMS_STUDY_NUMBER.isNotNull())\
.select('df2.*', 'status')
方法二:
q = """SELECT query_response_temp.*, IF( query_response_temp.CTMS_STUDY_NUMBER IS NOT NULL AND s3_parquet_file_temp.CTMS_STUDY_NUMBER IS NOT NULL, "existing", "new" ) as status FROM query_response_temp FULL OUTER JOIN s3_parquet_file_temp on query_response_temp.CTMS_STUDY_NUMBER = s3_parquet_file_temp.CTMS_STUDY_NUMBER AND query_response_temp.CTMS_SITE_NUMBER = s3_parquet_file_temp.CTMS_SITE_NUMBER AND query_response_temp.FULL_NAME_OF_PI = s3_parquet_file_temp.FULL_NAME_OF_PI AND query_response_temp.PRIMARY_ROLE_NAME = s3_parquet_file_temp.PRIMARY_ROLE_NAME AND query_response_temp.CENTER_NAME = s3_parquet_file_temp.CENTER_NAME"""
delta_records = sqlContext.sql(q)
delta_records = delta_records.filter(delta_records.CTMS_STUDY_NUMBER.isNotNull())
发生错误时发生错误的代码:
dataset = dataset.drop("hash_key")
print('hash_key dropped.')
dataset.coalesce(1).write.format('csv') \
.option('header', 'true') \
.option("compression", "none") \
.mode('overwrite') \
.save(_path)
【问题讨论】:
标签: python pandas amazon-s3 pyspark parquet