在RDD转换中调用类方法时出现PySpark的 CONTEXT_ONLY_VALID_ON_DRIVER错误(AWS Glue)
我在运行一个Glue作业,在调用存放在S3上的一个框架Python文件。
我按如下方式下载并导入框架:
os.system(f'aws s3 cp s3://s3_bucket/Common/ABC/XML_TO_PARQUE_FRAMEWORK.py ./ --region eu-west-1 --quiet')
import XML_TO_PARQUE_FRAMEWORK
XML_TO_PARQUE_FRAMEWORK.XmlToParquetFramework(
new_schema, detected_row_tag, submission_date, input_path, lg, config_dict, spark, glueContext,
mdrp_stage_bucket, trading_venue_code, file_id, file_id_length, AuroraConn, aurorahost, file_timestamp,
aurorauser, aurorapassword, auroraport, auroradb, mdrp_crawler_flag, stage_parquet_path,
curated_parquet_path, s3obj, s3
).run()
问题
在下面的方法中我遇到了错误:
def main_transaction_table(self, transaction_data_df, dynamodb_table):
batch_size = 100
df1 = transaction_data_df.select("tr_number", "aa_id").distinct()
df1.show(truncate=False)
result = df1.rdd.glom() \
.flatMap(lambda batch: [batch[i:i+batch_size] for i in range(0, len(batch), batch_size)]) \
.map(lambda batch: {
dynamodb_table: {
'Keys': [
{
'tr_number': {'S': str(row.tr_number)},
'aa_id': {'S': str(row.aa_id)}
} for row in batch
]
}
}) \
.flatMap(lambda req_item: self.get_duplicates(req_item, dynamodb_table)) \
.map(lambda item: Row(
tx_id=item['tr_number']['S'],
execgpty=item['aa_id']['S']
))
print('result is done')
schema2 = StructType([
StructField("tr_number", StringType(), False),
StructField("aa_id", StringType(), False)
])
print('schema 2 is done')
dup_df = self.spark.createDataFrame(result, schema2)
print('dup df done')
dup_df.persist()
错误
pyspark.errors.exceptions.base.PySparkRuntimeError: [CONTEXT_ONLY_VALID_ON_DRIVER]
It appears that you are attempting to reference SparkContext from a broadcast variable,
action, or transformation. SparkContext can only be used on the driver, not in code that
it run on workers.
观察
schema 2 is done会被打印出来,但在createDataFrame时发生错误。self.spark在类初始化期间被传递。self.get_duplicates()是同一类中的一个方法
我尝试过
- 将
get_duplicates变成一个@staticmethod - 将
get_duplicates移到类外,作为一个独立的函数 - 尝试将方法赋值给一个变量(
get_dups = self.get_duplicates) - 确保
get_duplicates不显式地使用spark
然而,上述方法都未能解决该错误。
解决方案
你的错误:
pyspark.errors.exceptions.base.PySparkRuntimeError: [CONTEXT_ONLY_VALID_ON_DRIVER]
It appears that you are attempting to reference SparkContext from a broadcast variable,
action, or transformation. SparkContext can only be used on the driver, not in code that
it run on workers.
含义:
当 rdd.map(), flatMap(), filter(),或在执行器上执行的类似转换时,Spark会对函数引用的所有对象进行序列化。如果该函数引用了self,并且类实例包含一个 SparkSession, SparkContext, DataFrame, or RDD,那么…… Spark会尝试序列化这些仅在驱动端存在的对象并失败。
解决方案:
class YourProcessor:
@staticmethod
def transform(x):
return x * 2
rdd.map(Processor.transform)
这避免了对仅在驱动端存在的Spark对象进行序列化,提升可维护性,并随着代码库的发展,防止出现类似的序列化错误。
站内所有文章版权归属LeftHeroAI导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。
