从Event Hub取数存AppendBlob后,Parquet读取报错求助
问题解析与解决方案
错误原因解析
1. OSError: Couldn't deserialize thrift: No more data to read. Deserializing page header failed.
Parquet文件依赖Thrift格式序列化元数据和数据页,这个错误说明文件结构已损坏:你将多个独立生成的完整Parquet文件的二进制内容直接追加到同一个AppendBlob中,导致Parquet解析器无法识别拼接后的无效结构,无法正常解析页头信息。
2. Unexpected end of stream: Page was smaller (4) than expected (13)
这个错误是文件损坏的衍生问题:读取Parquet数据页时,实际读取到的字节数远小于预期值,说明追加的二进制内容破坏了原有数据页的完整性,导致文件末尾被截断或结构错乱。
核心问题点
你的代码存在两个关键错误:
- Parquet文件不能直接二进制追加:Parquet是自包含的列式存储格式,每个
to_parquet()生成的是完整的Parquet文件,直接将多个完整文件的二进制内容拼接,会导致文件结构完全无效。 - BytesIO对象重复使用未重置:循环中复用同一个
BytesIO对象,每次写入后仅seek(0)但未清空,导致多次写入的Parquet内容叠加,进一步加剧文件损坏。
修正后的代码
import asyncio from datetime import datetime import pandas as pd from io import BytesIO from azure.storage.blob import BlobServiceClient from azure.eventhub.aio import EventHubConsumerClient from azure.eventhub.extensions.checkpointstoreblobaio import BlobCheckpointStore # 配置信息 EVENT_HUB_CONNECTION_STR = "" EVENT_HUB_NAME = "" BLOB_STORAGE_CONNECTION_STRING = "" BLOB_CONTAINER_NAME = "" # 提前初始化BlobServiceClient,避免重复创建 blob_service_client = BlobServiceClient.from_connection_string(BLOB_STORAGE_CONNECTION_STRING) async def on_event(partition_context, event): global finalDF try: data = event.body_as_json(encoding='UTF-8') df = pd.DataFrame(data, index=[0]) finalDF = pd.concat([finalDF, df]) if finalDF.shape[0] > 100: unique_bp_ids = finalDF['batteryserialnumber'].unique().tolist() current_datetime = datetime.now() year, month, day = current_datetime.year, current_datetime.month, current_datetime.day for bp_id in unique_bp_ids: temp_df = finalDF[finalDF['batteryserialnumber'] == bp_id] blob_path = f'new8_{year}/{month}/{bp_id}/{bp_id}_{year}_{month}_{day}.parquet' blob_client = blob_service_client.get_blob_client(container=BLOB_CONTAINER_NAME, blob=blob_path) # 处理已有Blob数据:如果存在则读取合并,否则直接写入 if blob_client.exists(): # 读取已有Parquet数据 existing_data = blob_client.download_blob().readall() existing_df = pd.read_parquet(BytesIO(existing_data)) # 合并新数据 combined_df = pd.concat([existing_df, temp_df]) else: combined_df = temp_df # 将合并后的数据写入Parquet并上传(使用BlockBlob,覆盖原有文件) parquet_buffer = BytesIO() combined_df.to_parquet(parquet_buffer) parquet_buffer.seek(0) blob_client.upload_blob(data=parquet_buffer, overwrite=True, blob_type='BlockBlob') finalDF = pd.DataFrame() print('数据处理完成并上传') except Exception as e: print(f'错误发生: {e}') await partition_context.update_checkpoint(event) async def main(): checkpoint_store = BlobCheckpointStore.from_connection_string( BLOB_STORAGE_CONNECTION_STRING, BLOB_CONTAINER_NAME ) client = EventHubConsumerClient.from_connection_string( EVENT_HUB_CONNECTION_STR, consumer_group="$Default", checkpoint_store=checkpoint_store, eventhub_name=EVENT_HUB_NAME, ) async with client: await client.receive(on_event=on_event, starting_position="-1") if __name__ == "__main__": finalDF = pd.DataFrame() loop = asyncio.get_event_loop() loop.run_until_complete(main())
关键修改说明
- 提前初始化BlobServiceClient:避免在循环中重复创建客户端,提升性能。
- 合并数据而非二进制追加:读取已有Blob中的Parquet数据,与新数据合并后重新生成完整的Parquet文件上传,保证文件结构有效。
- 使用BlockBlob替代AppendBlob:因为Parquet文件需要完整的结构,BlockBlob支持覆盖写入,更适合这种场景。
- 循环内独立创建BytesIO:每次处理一个电池序列号时都新建BytesIO对象,避免数据叠加。
内容的提问来源于stack exchange,提问作者zircon
相关产品推荐
相关产品推荐

