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
相关产品推荐
相关产品推荐

