如何通过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
相关产品推荐
相关产品推荐

