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

Python如何拆分Azure Service Bus批量字符串消息并归档到Azure存储

方案可行性结论

你设计的方案完全可行,逻辑链路清晰,针对单条消息包含50个拼接JSON的量级,性能和稳定性都完全满足需求,没有明显缺陷。

具体实现建议

1. 无分隔符拼接JSON的正确拆分

不要用正则匹配大括号的方式拆分,容易因为JSON内部嵌套大括号、字符串中包含大括号的场景出错,推荐直接用Python标准库json模块自带的JSONDecoder.raw_decode方法,它可以从字符串的起始位置解析出一个完整JSON对象,同时返回该对象结束的下标,循环调用即可拆分所有拼接的JSON。
拆分函数示例代码:

import json
from typing import List

def split_concat_json(raw_str: str) -> List[dict]:
    decoder = json.JSONDecoder()
    result = []
    idx = 0
    raw_str = raw_str.strip()
    while idx < len(raw_str):
        # 跳过空白字符
        while idx < len(raw_str) and raw_str[idx].isspace():
            idx += 1
        if idx >= len(raw_str):
            break
        # 解析单个JSON,返回解析结果和结束下标
        obj, end_idx = decoder.raw_decode(raw_str, idx)
        result.append(obj)
        idx = end_idx
    return result

2. 构造DataFrame并写入Blob存储

构造DataFrame

拆分得到的字典列表可以直接传入pandas生成DataFrame,不需要额外处理字段:

import pandas as pd

json_list = split_concat_json(str(msg))
df = pd.DataFrame(json_list)

写入Azure Blob存储

如果是追加归档,优先选择追加Blob类型,适合增量写入的场景,也可以按时间维度(比如按天/按小时)生成新的Blob文件,避免单个文件过大。
需要先安装依赖:pip install azure-storage-blob pandas
写入示例代码:

from azure.storage.blob import BlobServiceClient, BlobType

# 初始化Blob客户端
blob_conn_str = "你的Blob存储连接字符串"
container_name = "归档容器名"
# 按日期生成归档文件名,方便后续检索
blob_name = f"service_bus_archive/20240520.csv"

blob_service_client = BlobServiceClient.from_connection_string(blob_conn_str)
blob_client = blob_service_client.get_blob_client(container=container_name, blob=blob_name)

# 把DataFrame转成CSV格式字符串
csv_content = df.to_csv(index=False, header=not blob_client.exists())

# 追加写入
if blob_client.exists():
    blob_client.append_block(csv_content)
else:
    # 第一次写入要先创建追加Blob
    blob_client.upload_blob(csv_content, blob_type=BlobType.APPENDBLOB)

3. 原有Service Bus读取代码的优化

你原有的代码可以直接整合上面的逻辑,同时补充异常处理,避免解析失败时丢失数据:

from azure.servicebus import ServiceBusClient
import json
import pandas as pd
from azure.storage.blob import BlobServiceClient, BlobType
from typing import List
from datetime import datetime

# 配置项
sb_conn_str = "**"
topic_name = "***"
subscription_name = "***"
blob_conn_str = "你的Blob存储连接字符串"
container_name = "归档容器名"

def split_concat_json(raw_str: str) -> List[dict]:
    decoder = json.JSONDecoder()
    result = []
    idx = 0
    raw_str = raw_str.strip()
    while idx < len(raw_str):
        while idx < len(raw_str) and raw_str[idx].isspace():
            idx += 1
        if idx >= len(raw_str):
            break
        obj, end_idx = decoder.raw_decode(raw_str, idx)
        result.append(obj)
        idx = end_idx
    return result

def archive_to_blob(df: pd.DataFrame):
    # 按天生成归档文件
    today = datetime.utcnow().strftime("%Y%m%d")
    blob_name = f"service_bus_archive/{today}.csv"
    blob_service_client = BlobServiceClient.from_connection_string(blob_conn_str)
    blob_client = blob_service_client.get_blob_client(container=container_name, blob=blob_name)
    csv_content = df.to_csv(index=False, header=not blob_client.exists())
    if blob_client.exists():
        blob_client.append_block(csv_content)
    else:
        blob_client.upload_blob(csv_content, blob_type=BlobType.APPENDBLOB)

servicebus_client = ServiceBusClient.from_connection_string(
    conn_str=sb_conn_str, logging_enable=True)

with servicebus_client:
    receiver = servicebus_client.get_subscription_receiver(
        topic_name=topic_name, subscription_name=subscription_name)
    with receiver:
        for msg in receiver:
            try:
                raw_content = str(msg)
                json_list = split_concat_json(raw_content)
                df = pd.DataFrame(json_list)
                archive_to_blob(df)
                receiver.complete_message(msg)
            except Exception as e:
                # 解析失败时将消息放回队列或者移入死信队列,避免丢数据
                print(f"处理消息失败:{str(e)}")
                receiver.abandon_message(msg)
额外注意事项
  • 如果你使用的是Azure函数,不需要自己维护Service Bus客户端的生命周期,可以直接用Azure函数的Service Bus触发器,比时间触发更高效,消息到达就会触发执行,不需要轮询
  • 如果归档数据不需要频繁读取,可以在写入Blob时设置归档存储层,进一步降低存储成本
  • 单条追加Blob的最大大小是4.75TB,足够日常归档使用,如果数据量极大可以按小时拆分文件

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.29 07:36:05