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

从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())

关键修改说明

  1. 提前初始化BlobServiceClient:避免在循环中重复创建客户端,提升性能。
  2. 合并数据而非二进制追加:读取已有Blob中的Parquet数据,与新数据合并后重新生成完整的Parquet文件上传,保证文件结构有效。
  3. 使用BlockBlob替代AppendBlob:因为Parquet文件需要完整的结构,BlockBlob支持覆盖写入,更适合这种场景。
  4. 循环内独立创建BytesIO:每次处理一个电池序列号时都新建BytesIO对象,避免数据叠加。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.02 10:58:12