在RDD转换中调用类方法时出现PySpark的 CONTEXT_ONLY_VALID_ON_DRIVER错误(AWS Glue)

编程语言 2026-07-08

我在运行一个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.

观察

  1. schema 2 is done 会被打印出来,但在 createDataFrame 时发生错误。
  2. self.spark 在类初始化期间被传递。
  3. self.get_duplicates() 是同一类中的一个方法

我尝试过

  1. get_duplicates 变成一个 @staticmethod
  2. get_duplicates 移到类外,作为一个独立的函数
  3. 尝试将方法赋值给一个变量(get_dups = self.get_duplicates
  4. 确保 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导航站,无授权禁止任何主体转载、抄袭、复制内容,亦不得私自架设镜像站点。一经侵权,本站将通过法律途径追责。

相关文章