AWS Glue将PySpark DF导入Redshift:含空值列的处理方案求助
可行替代方案
1. 升级Redshift JDBC驱动版本
已知v2.1.0.9版本存在该空值处理bug,在后续发布的v2.1.1.0及以上版本中已修复。操作步骤:
- 下载官方最新版Redshift JDBC驱动(如v2.1.17.1001)
- 将驱动包上传至S3存储桶
- 在AWS Glue Job配置中,通过
--extra-jars参数指定该驱动包的S3路径,覆盖默认旧版本驱动
2. 精准匹配目标列类型转换空值列
之前转换NullType到StringType仍报错,大概率是因为转换后列仍全为空,驱动无法识别。需针对Redshift目标表的列类型逐一转换,确保类型完全匹配:
from pyspark.sql.types import StringType, IntegerType, DateType from pyspark.sql.functions import col # 根据Redshift目标列类型转换 # 示例:目标列"desc"为VARCHAR类型,原列为NullType df = df.withColumn("desc", col("desc").cast(StringType())) # 示例:目标列"age"为INT类型,原列为NullType df = df.withColumn("age", col("age").cast(IntegerType())) # 示例:目标列"create_date"为DATE类型,原列为NullType df = df.withColumn("create_date", col("create_date").cast(DateType()))
转换后直接写入Redshift,原生null会被正确识别,无需额外替换步骤。
3. 使用AWS Glue专用Redshift连接器(推荐)
放弃JDBC驱动,改用Glue内置的Redshift连接器——它基于Redshift COPY命令,性能更优且能自动处理空值问题:
# 转换为DynamicFrame(可选,也可直接用DataFrame) dynamic_frame = glueContext.create_dynamic_frame.from_df(df, name="mongo_data") # 写入Redshift glueContext.write_dynamic_frame.from_jdbc_conf( frame=dynamic_frame, catalog_connection="your_redshift_connection", # 提前在Glue Data Catalog配置的连接 connection_options={ "dbtable": "target_schema.target_table", "database": "your_redshift_db" }, redshift_tmp_dir="s3://your-bucket/temp_redshift/", # 临时存储路径,需配置IAM权限 transformation_ctx="write_redshift" )
也可直接用DataFrame写入:
df.write \ .format("io.github.spark_redshift_community.spark.redshift") \ .option("url", "jdbc:redshift://cluster-endpoint:5439/db_name?user=username&password=password") \ .option("dbtable", "target_schema.target_table") \ .option("tempdir", "s3://your-bucket/temp/") \ .option("aws_iam_role", "arn:aws:iam::your-account-id:role/redshift-access-role") \ .mode("append") \ .save()
4. 调整JDBC写入参数适配空值
如果暂时无法升级驱动或切换连接器,尝试设置JDBC的nullValue参数,针对不同列类型指定合法空值:
- 字符串类型列,设置
nullValue为空字符串:df.write \ .format("jdbc") \ .option("url", "jdbc:redshift://cluster-url:5439/db") \ .option("dbtable", "target_table") \ .option("user", "db_user") \ .option("password", "db_pass") \ .option("nullValue", "") \ .mode("append") \ .save() - 数值类型列,设置
nullValue为NULL(Redshift会自动解析为原生空值):df.write \ .format("jdbc") \ .option("url", "jdbc:redshift://cluster-url:5439/db") \ .option("dbtable", "target_table") \ .option("user", "db_user") \ .option("password", "db_pass") \ .option("nullValue", "NULL") \ .mode("append") \ .save()
内容的提问来源于stack exchange,提问作者Sanket Kelkar
相关产品推荐
相关产品推荐

