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

如何通过Dask连接存储在MinIO上的Delta Lake?

问题描述

我可以通过deltalake Python包直接读取存储在MinIO上的Delta Lake表,代码如下:

storage_options = {
    "AWS_ENDPOINT_URL": "http://localhost:9000",
    "AWS_REGION": "local",
    "AWS_ACCESS_KEY_ID": access_key,
    "AWS_SECRET_ACCESS_KEY": secret_key,
    "AWS_S3_ALLOW_UNSAFE_RENAME": "true",
    "AWS_ALLOW_HTTP": "true"
}

dt = DeltaTable("s3a://my_bucket/data", storage_options=storage_options)
df = dt.to_pandas()

但当我尝试使用dask-deltatable读取该表并转换为Dask DataFrame时,程序仍试图连接AWS服务,报错信息如下:

OSError                                   Traceback (most recent call last)
Cell In[3], line 1
----> 1 ddf = dask_deltatable.read_deltalake("s3a://my_bucket/data", storage_options=storage_options)

File ~/.local/lib/python3.10/site-packages/dask_deltatable/core.py:285, in read_deltalake(path, catalog, database_name, table_name, version, columns, storage_options, datetime, delta_storage_options, **kwargs)
    282         raise ValueError("Please Provide Delta Table path")
    284     delta_storage_options = utils.maybe_set_aws_credentials(path, delta_storage_options)  # type: ignore
---> 285     resultdf = _read_from_filesystem(
    286         path=path,
    287         version=version,
    288         columns=columns,
    289         storage_options=storage_options,
    290         datetime=datetime,
    291         delta_storage_options=delta_storage_options,
    292         **kwargs,
    293     )
    294 return resultdf

File ~/.local/lib/python3.10/site-packages/dask_deltatable/core.py:102, in _read_from_filesystem(path, version, columns, datetime, storage_options, delta_storage_options, **kwargs)
     99 delta_storage_options = utils.maybe_set_aws_credentials(path, delta_storage_options)  # type: ignore
    101 fs, fs_token, _ = get_fs_token_paths(path, storage_options=storage_options)
---> 102 dt = DeltaTable(
    103     table_uri=path, version=version, storage_options=delta_storage_options
    104 )
    105 if datetime is not None:
    106     dt.load_as_version(datetime)

File ~/.local/lib/python3.10/site-packages/deltalake/table.py:297, in DeltaTable.__init__(self, table_uri, version, storage_options, without_files, log_buffer_size)
    277 """
    278 Create the Delta Table from a path with an optional version.
    279 Multiple StorageBackends are currently supported: AWS S3, Azure Data Lake Storage Gen2, Google Cloud Storage (GCS) and local URI.
   (...)
    294 
    295 """
    296 self._storage_options = storage_options
---> 297 self._table = RawDeltaTable(
    298     str(table_uri),
    299     version=version,
    300     storage_options=storage_options,
    301     without_files=without_files,
    302     log_buffer_size=log_buffer_size,
    303 )

OSError: Generic S3 error: Error after 10 retries in 13.6945151s, max_retries:10, retry_timeout:180s, source:error sending request for url (http://169.254.169.254/latest/api/token)

请问如何成功实现从MinIO读取Delta Lake表到Dask DataFrame?

解决方案

问题根源在于dask-deltatable内置的utils.maybe_set_aws_credentials函数会自动处理AWS凭据,覆盖MinIO的配置,同时触发程序尝试从AWS元数据服务(即报错中的169.254.169.254地址)获取凭据,导致连接失败。以下是三种可行的解决方法:

方法一:通过delta_storage_options传递MinIO配置

dask-deltatable的read_deltalake函数提供了delta_storage_options参数,专门用于传递给底层的DeltaTable实例。直接将MinIO的配置放在该参数中,避免被自动处理逻辑干扰:

delta_storage_options = {
    "AWS_ENDPOINT_URL": "http://localhost:9000",
    "AWS_REGION": "local",
    "AWS_ACCESS_KEY_ID": access_key,
    "AWS_SECRET_ACCESS_KEY": secret_key,
    "AWS_S3_ALLOW_UNSAFE_RENAME": "true",
    "AWS_ALLOW_HTTP": "true",
    "AWS_DISABLE_SSL": "true"  # HTTP连接时建议添加,防止SSL校验问题
}

ddf = dask_deltatable.read_deltalake(
    "s3a://my_bucket/data",
    delta_storage_options=delta_storage_options
)

方法二:禁用AWS元数据服务访问

如果需要保留storage_options参数,可添加AWS_EC2_METADATA_DISABLED: "true"阻止程序访问AWS元数据服务,同时确保delta_storage_options也包含完整的MinIO配置:

storage_options = {
    "AWS_ENDPOINT_URL": "http://localhost:9000",
    "AWS_REGION": "local",
    "AWS_ACCESS_KEY_ID": access_key,
    "AWS_SECRET_ACCESS_KEY": secret_key,
    "AWS_S3_ALLOW_UNSAFE_RENAME": "true",
    "AWS_ALLOW_HTTP": "true",
    "AWS_EC2_METADATA_DISABLED": "true"
}

ddf = dask_deltatable.read_deltalake(
    "s3a://my_bucket/data",
    storage_options=storage_options,
    delta_storage_options=storage_options
)

方法三:手动加载DeltaTable再转换为Dask DataFrame

绕过dask-deltatable的自动凭据处理逻辑,先用deltalake包加载MinIO上的Delta表,再通过dask_deltatable.from_delta_table方法转换为Dask DataFrame:

from deltalake import DeltaTable
import dask_deltatable

storage_options = {
    "AWS_ENDPOINT_URL": "http://localhost:9000",
    "AWS_REGION": "local",
    "AWS_ACCESS_KEY_ID": access_key,
    "AWS_SECRET_ACCESS_KEY": secret_key,
    "AWS_S3_ALLOW_UNSAFE_RENAME": "true",
    "AWS_ALLOW_HTTP": "true"
}

# 先用deltalake加载表
dt = DeltaTable("s3a://my_bucket/data", storage_options=storage_options)
# 转换为Dask DataFrame
ddf = dask_deltatable.from_delta_table(dt)

内容的提问来源于stack exchange,提问作者James Baker

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 17:54:52