You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在Databricks读取S3 Delta表遇错,如何确保生成_delta_log?

问题描述

在Databricks中读取S3存储桶内的Delta表时,小表(如table_1)可正常加载,但被AWS Glue判定为“大表”的table_2会抛出如下错误:

AnalysisException: Incompatible format detected.
You are trying to read from SECOND_DATALAKE/table_2/ using Delta, but there is no transaction log present. Check the upstream job to make sure that it is writing using format("delta") and that you are trying to read from the table base path.

推测问题根源是AWS Glue作业处理大表时未生成_delta_log文件夹,以下是Glue作业的代码片段:

table_to_copy = 'TABLE_2'
try:
    # Script generado para node Oracle SQL
    OracleSQL_node1 = glueContext.create_dynamic_frame.from_options(
        connection_type="oracle",
        connection_options={
                "url": 'jdbc:oracle:thin://@datalake.eu-west-3.rds.amazonaws.com:10523:NAME1', 
                "user": username,
                "password": password,
                "dbtable": 'SCHEMA.' + table_to_copy,
            },
        transformation_ctx="OracleSQL_node1",
    )
    
    # Convertir a DataFrame
    dataFrame = OracleSQL_node1.toDF()
    
    # Escribir en formato Delta
    dataFrame.repartition(10) \
        .write.format('delta') \
        .mode('overwrite') \
        .save("SECOND_DATALAKE" + table_to_copy)  
        
    logger.info(f"{table_to_copy} was copied to the Delta Lake correctly")
except Exception as e:
    logger.error(f"Error copying {table_to_copy} to Delta Lake: {str(e)}")

job.commit()
spark.stop()
解决方案与优化建议

1. 确保Delta Lake依赖正确加载

AWS Glue默认未内置完整的Delta Lake依赖,处理大表时需显式指定加载Delta格式:

  • 在Glue作业的作业参数中添加--datalake-formats delta(适用于Glue 3.0及以上版本),该参数会自动加载匹配Glue版本的Delta Lake依赖包,避免版本不兼容导致的日志生成失败。

2. 修复写入路径格式错误

原代码中路径拼接缺少斜杠,可能导致数据写入到错误的S3路径,进而无法找到_delta_log:
将写入路径修改为:

.save(f"SECOND_DATALAKE/{table_to_copy}")

确保路径为SECOND_DATALAKE/table_2/格式,而非拼接后无分隔符的错误路径。

3. 移除提前终止Spark会话的代码

原代码中job.commit()后调用spark.stop(),可能在Delta事务日志未完全写入时终止Spark进程,导致_delta_log生成不完整。Glue作业完成后会自动管理Spark会话生命周期,因此直接删除spark.stop()语句即可。

4. 优化大表写入的Delta配置

添加显式的Delta配置项,确保事务日志稳定生成:

dataFrame.repartition(200)  # 根据数据量调整分区数,建议每分区1-2GB
    .write.format('delta')
    .mode('overwrite')
    .option("delta.logRetentionDuration", "interval 30 days")
    .option("delta.optimizeWrite", "true")  # 自动优化写入布局,提升大表写入稳定性
    .save(f"SECOND_DATALAKE/{table_to_copy}")

同时可在Glue作业的Spark配置中添加:

--conf spark.databricks.delta.commitInfo.enabled=true

强制Delta记录完整的提交信息,确保日志生成。

5. 改用Glue DynamicFrame直接写入Delta

避免转换为Spark DataFrame的开销,利用Glue原生优化提升大表写入稳定性:

glueContext.write_dynamic_frame.from_options(
    frame=OracleSQL_node1,
    connection_type="s3",
    connection_options={"path": f"SECOND_DATALAKE/{table_to_copy}"},
    format="delta",
    transformation_ctx="write_delta_node"
)

这种方式无需手动转换DataFrame,Glue会自动处理大表的分区与写入逻辑,降低日志生成失败的概率。

6. 增强写入后的校验逻辑

添加S3路径校验,确保_delta_log生成后再标记作业成功:

import boto3

# 提取S3桶名和前缀(需根据实际SECOND_DATALAKE的路径修改,假设SECOND_DATALAKE是s3://bucket-name/前缀)
bucket_name = "your-s3-bucket-name"
delta_log_prefix = f"{table_to_copy}/_delta_log/"

s3_client = boto3.client('s3')
response = s3_client.list_objects_v2(Bucket=bucket_name, Prefix=delta_log_prefix)

if 'Contents' not in response:
    raise Exception(f"Delta transaction log not found for table {table_to_copy}, write failed")

将这段代码放在写入操作之后、日志打印之前,避免假成功的情况。

修改后的完整代码示例
import boto3

table_to_copy = 'TABLE_2'
try:
    # 读取Oracle数据
    OracleSQL_node1 = glueContext.create_dynamic_frame.from_options(
        connection_type="oracle",
        connection_options={
                "url": 'jdbc:oracle:thin://@datalake.eu-west-3.rds.amazonaws.com:10523:NAME1', 
                "user": username,
                "password": password,
                "dbtable": 'SCHEMA.' + table_to_copy,
            },
        transformation_ctx="OracleSQL_node1",
    )
    
    # 直接用DynamicFrame写入Delta
    glueContext.write_dynamic_frame.from_options(
        frame=OracleSQL_node1,
        connection_type="s3",
        connection_options={"path": f"SECOND_DATALAKE/{table_to_copy}"},
        format="delta",
        transformation_ctx="write_delta_node"
    )
    
    # 校验_delta_log是否存在
    bucket_name = "your-s3-bucket-name"
    delta_log_prefix = f"{table_to_copy}/_delta_log/"
    s3_client = boto3.client('s3')
    response = s3_client.list_objects_v2(Bucket=bucket_name, Prefix=delta_log_prefix)
    
    if 'Contents' not in response:
        raise Exception(f"Delta transaction log missing for {table_to_copy}")
    
    logger.info(f"{table_to_copy} was copied to Delta Lake successfully")
except Exception as e:
    logger.error(f"Error copying {table_to_copy} to Delta Lake: {str(e)}")
    raise  # 抛出异常标记作业失败

job.commit()

内容的提问来源于stack exchange,提问作者jhonny gonzalez

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.04 11:07:46