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]" # }
相关产品推荐
相关产品推荐

