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

Apache Spark 3.4中toPandas().to_csv()认证失败问题求助

问题描述

我们使用Azure Synapse,微软通知Apache Spark 3.3版本将于2025年停用。为提前准备,创建了Spark 3.4版本新池并迁移笔记本,其他功能正常,但使用toPandas().to_csv()向Gen2数据湖写入单个CSV时出现认证错误,移回3.3池则正常。

已确认工作区拥有Storage Blob Data Contributor权限,配置了存储目录链接服务器并添加storage_options: { linked_server: '[linkedserver]'}。推测问题源于包需更新或存在bug,寻求原因及解决方案。

相关代码

localPathForFiles = "abfss://[container]@[storageaccount].dfs.core.windows.net/extracts/customer-data/{0}/".format(CustomerName)
Df.toPandas().to_csv(localPathForFiles + TaskTypesFileName, lineterminator="\r\n", index=False, header=True, quotechar='"', quoting=csv.QUOTE_NONNUMERIC)

完整错误信息

ClientAuthenticationError                 Traceback (most recent call last)
File ~/cluster-env/clonedenv/lib/python3.10/site-packages/adlfs/spec.py:2030, in AzureBlobFile._async_upload_chunk(self, final, **kwargs)
2027 async with self.container_client.get_blob_client(
2028     blob=self.blob
2029 ) as bc:
2030     await bc.stage_block(
2031         block_id=block_id,
2032         data=chunk,
2033         length=len(chunk),
2034     )
2035     self._block_list.append(block_id)

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/azure/core/tracing/decorator_async.py:77, in distributed_trace_async.<locals>.decorator.<locals>.wrapper_use_tracer(*args, **kwargs)
76 if span_impl_type is None:
77     return await func(*args, **kwargs)
79 # Merge span is parameter is set, but only if no explicit parent are passed

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/azure/storage/blob/aio/_blob_client_async.py:1644, in BlobClient.stage_block(self, block_id, data, length, **kwargs)
1643 except HttpResponseError as error:
1644     process_storage_error(error)

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/azure/storage/blob/_shared/response_handlers.py:184, in process_storage_error(storage_error)
182 try:
183     # `from None` prevents us from double printing the exception (suppresses generated layer error context)
184     exec("raise error from None")   # pylint: disable=exec-used # nosec
185 except SyntaxError as exc:

File <string>:1

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/azure/storage/blob/aio/_blob_client_async.py:1642, in BlobClient.stage_block(self, block_id, data, length, **kwargs)
1641 try:
1642     return await self._client.block_blob.stage_block(**options)
1643 except HttpResponseError as error:

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/azure/core/tracing/decorator_async.py:77, in distributed_trace_async.<locals>.decorator.<locals>.wrapper_use_tracer(*args, **kwargs)
76 if span_impl_type is None:
77     return await func(*args, **kwargs)
79 # Merge span is parameter is set, but only if no explicit parent are passed

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/azure/storage/blob/_generated/aio/operations/_block_blob_operations.py:645, in BlockBlobOperations.stage_block(self, block_id, content_length, body, transactional_content_md5, transactional_content_crc64, timeout, request_id_parameter, lease_access_conditions, cpk_info, cpk_scope_info, **kwargs)
644 if response.status_code not in [201]:
645     map_error(status_code=response.status_code, response=response, error_map=error_map)
646     error = self._deserialize.failsafe_deserialize(_models.StorageError, pipeline_response)

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/azure/core/exceptions.py:164, in map_error(status_code, response, error_map)
163 error = error_type(response=response)
164 raise error

ClientAuthenticationError: Server failed to authenticate the request. Please refer to the information in the www-authenticate header.
RequestId:b7074898-501e-007f-4337-e989a8000000
Time:2024-08-08T02:07:05.0273937Z
ErrorCode:NoAuthenticationInformation
Content: <?xml version="1.0" encoding="utf-8"?><Error><Code>NoAuthenticationInformation</Code><Message>Server failed to authenticate the request. Please refer to the information in the www-authenticate header.
RequestId:b7074898-501e-007f-4337-e989a8000000
Time:2024-08-08T02:07:05.0273937Z</Message></Error>

The above exception was the direct cause of the following exception:

RuntimeError                              Traceback (most recent call last)
Cell In[15], line 9
7 typesDf = spark.read.jdbc(url=serverCommsConnectionString, table=typesSql, properties=connectionCommsProperties)
8 #Write to storage
9 typesDf.toPandas().to_csv(localPathForFiles + TaskTypesFileName, lineterminator="\r\n", index=False, header=True, quotechar='"', quoting=csv.QUOTE_NONNUMERIC)        

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/pandas/util/_decorators.py:211, in deprecate_kwarg.<locals>._deprecate_kwarg.<locals>.wrapper(*args, **kwargs)
209     else:
210         kwargs[new_arg_name] = new_arg_value
211 return func(*args, **kwargs)

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/pandas/core/generic.py:3720, in NDFrame.to_csv(self, path_or_buf, sep, na_rep, float_format, columns, header, index, index_label, mode, encoding, compression, quoting, quotechar, lineterminator, chunksize, date_format, doublequote, escapechar, decimal, errors, storage_options)
3709 df = self if isinstance(self, ABCDataFrame) else self.to_frame()
3711 formatter = DataFrameFormatter(
3712     frame=df,
3713     header=header,
>    (...)
3717     decimal=decimal,
3718 )
3720 return DataFrameRenderer(formatter).to_csv(
3721     path_or_buf,
3722     lineterminator=lineterminator,
3723     sep=sep,
3724     encoding=encoding,
3725     errors=errors,
3726     compression=compression,
3727     quoting=quoting,
3728     columns=columns,
3729     index_label=index_label,
3730     mode=mode,
3731     chunksize=chunksize,
3732     quotechar=quotechar,
3733     date_format=date_format,
3734     doublequote=doublequote,
3735     escapechar=escapechar,
3736     storage_options=storage_options,
3737 )

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/pandas/util/_decorators.py:211, in deprecate_kwarg.<locals>._deprecate_kwarg.<locals>.wrapper(*args, **kwargs)
209     else:
210         kwargs[new_arg_name] = new_arg_value
211 return func(*args, **kwargs)

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/pandas/io/formats/format.py:1189, in DataFrameRenderer.to_csv(self, path_or_buf, encoding, sep, columns, index_label, mode, compression, quoting, quotechar, lineterminator, chunksize, date_format, doublequote, escapechar, errors, storage_options)
1168     created_buffer = False
1170 csv_formatter = CSVFormatter(
1171     path_or_buf=path_or_buf,
1172     lineterminator=lineterminator,
>    (...)
1187     formatter=self.fmt,
1188 )
1189 csv_formatter.save()
1191 if created_buffer:
1192     assert isinstance(path_or_buf, StringIO)

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/pandas/io/formats/csvs.py:241, in CSVFormatter.save(self)
237 """
238 Create the writer & save.
239 """
240 # apply compression and byte/text conversion
241 with get_handle(
242     self.filepath_or_buffer,
243     self.mode,
244     encoding=self.encoding,
245     errors=self.errors,
246     compression=self.compression,
247     storage_options=self.storage_options,
248 ) as handles:
249 
250     # Note: self.encoding is irrelevant here
251     self.writer = csvlib.writer(
252         handles.handle,
253         lineterminator=self.lineterminator,
>    (...)
258         quotechar=self.quotechar,
259     )
261     self._save()

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/pandas/io/common.py:133, in IOHandles.__exit__(self, *args)
132 def __exit__(self, *args: Any) -> None:
133     self.close()

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/pandas/io/common.py:125, in IOHandles.close(self)
123     self.created_handles.remove(self.handle)
124 for handle in self.created_handles:
125     handle.close()
126 self.created_handles = []
127 self.is_wrapped = False

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/adlfs/spec.py:1903, in AzureBlobFile.close(self)
1901 """Close file and azure client."""
1902 asyncio.run_coroutine_threadsafe(close_container_client(self), loop=self.loop)
1903 super().close()

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/fsspec/spec.py:1932, in AbstractBufferedFile.close(self)
1930 else:
1931     if not self.forced:
1932         self.flush(force=True)
1934     if self.fs is not None:
1935         self.fs.invalidate_cache(self.path)

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/fsspec/spec.py:1803, in AbstractBufferedFile.flush(self, force)
1800         self.closed = True
1801         raise
1803 if self._upload_chunk(final=force) is not False:
1804     self.offset += self.buffer.seek(0, 2)
1805     self.buffer = io.BytesIO()

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/fsspec/asyn.py:118, in sync_wrapper.<locals>.wrapper(*args, **kwargs)
115 @functools.wraps(func)
116 def wrapper(*args, **kwargs):
117     self = obj or args[0]
118     return sync(self.loop, func, *args, **kwargs)

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/fsspec/asyn.py:103, in sync(loop, func, timeout, *args, **kwargs)
101     raise FSTimeoutError from return_result
102 elif isinstance(return_result, BaseException):
103     raise return_result
104 else:
105     return return_result

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/fsspec/asyn.py:56, in _runner(event, coro, result, timeout)
54     coro = asyncio.wait_for(coro, timeout=timeout)
55 try:
56     result[0] = await coro
57 except Exception as ex:
58     result[0] = ex

File ~/cluster-env/clonedenv/lib/python3.10/site-packages/adlfs/spec.py:2067, in AzureBlobFile._async_upload_chunk(self, final, **kwargs)
2063                 await bc.commit_block_list(
2064                     block_list=block_list, metadata=self.metadata
2065                 )
2066         else:
2067             raise RuntimeError(f"Failed to upload block{e}!") from e
2068 elif self.mode == "ab":
2069     async with self.container_client.get_blob_client(blob=self.blob) as bc:

RuntimeError: Failed to upload blockServer failed to authenticate the request. Please refer to the information in the www-authenticate header.
RequestId:b7074898-501e-007f-4337-e989a8000000
Time:2024-08-08T02:07:05.0273937Z
ErrorCode:NoAuthenticationInformation
Content: <?xml version="1.0" encoding="utf-8"?><Error><Code>NoAuthenticationInformation</Code><Message>Server failed to authenticate the request. Please refer to the information in the www-authenticate header.
RequestId:b7074898-501e-007f-4337-e989a8000000
Time:2024-08-08T02:07:05.0273937Z</Message></Error>!
原因分析

这个问题的核心是Spark 3.4环境中,pandas.to_csv()依赖的adlfs或Azure存储相关包的认证逻辑与Spark 3.3存在差异:

  • Spark 3.4默认的Python 3.10环境搭配的adlfs版本,存在认证适配问题,无法正确使用Synapse工作区的托管身份或链接服务器配置完成身份验证。
  • 当通过toPandas()转换为Pandas DataFrame后,后续的to_csv()操作不再走Spark的存储认证逻辑,而是直接调用adlfs库访问ADLS Gen2,而该库在Spark 3.4环境中的配置未正确继承Synapse的身份信息。
解决方案

方案1:使用Spark原生API替代Pandas写入

放弃toPandas().to_csv()的方式,改用Spark原生写入API,复用Spark已配置的存储认证逻辑,避免Pandas层面的认证问题:

# 合并分区后写入单个CSV文件
Df.coalesce(1) \
  .write \
  .format("csv") \
  .option("header", "true") \
  .option("quote", "\"") \
  .option("quoteMode", "NON_NUMERIC") \
  .option("lineSep", "\r\n") \
  .mode("overwrite") \
  .save(localPathForFiles + "temp-" + TaskTypesFileName)

# 重命名Spark生成的临时文件为目标文件名
import os
from pyspark.sql import functions as F

# 获取生成的CSV文件路径
csv_files = dbutils.fs.ls(localPathForFiles + "temp-" + TaskTypesFileName)
csv_file = [f.path for f in csv_files if f.path.endswith(".csv")][0]

# 移动并重命名文件
dbutils.fs.mv(csv_file, localPathForFiles + TaskTypesFileName)
# 删除临时目录
dbutils.fs.rm(localPathForFiles + "temp-" + TaskTypesFileName, recurse=True)

方案2:显式传递认证信息给Pandas的storage_options

如果必须使用Pandas的to_csv(),调用时显式传入正确的认证参数,确保adlfs能获取有效身份凭证:

# 使用Synapse链接服务器配置
storage_options = {
    "linked_service": "[linkedserver]"
}

# 或直接指定存储账户密钥(若业务允许)
# storage_options = {
#     "account_key": "[your-storage-account-key]"
# }
相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 00:39:40