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

Azure函数中实现捕获多类异常并保留原消息的装饰器方案咨询

Azure函数中实现捕获多类异常并保留原消息的装饰器方案咨询

我完全理解你的需求——你需要一个能包裹Azure函数主逻辑的装饰器,既能捕获所有可能的异常(不管是JSON解析错误、数据校验错误还是接口调用的临时错误),又能牢牢保留原始的Event Hub消息,同时把异常详情打包成新事件发送出去。之前你尝试写装饰器时丢失了原消息,问题出在装饰器没有正确接收并绑定函数的输入参数,下面给你一套完整的可落地解决方案:

核心思路

我们要实现一个异步兼容的装饰器(因为你的main函数是异步函数),它会:

  1. 直接接收原始的eventHubMessage参数
  2. 调用你的业务逻辑函数并捕获所有异常
  3. 将原消息、异常类型、异常详情甚至堆栈信息打包成结构化的错误事件
  4. 统一发送错误事件到指定的异常处理通道(比如另一个Event Hub)

完整代码实现

以下是改造后的完整代码,包含装饰器实现、错误事件发送逻辑,以及简化后的业务主函数:

from fastapi import HTTPException
from pydantic import BaseModel
from dam_common.common import settings_helper
from enum import Enum
import json
import httpx
import logging
import functools
import traceback
from datetime import datetime
# 导入Azure Event Hub SDK用于发送错误事件
from azure.eventhub import AsyncEventHubProducerClient, EventData


class EventHubRequest(BaseModel):
    IHDRequestID: str
    ProductType: str
    RequestType: str
    RequesterADGroupName: str
    TimeStamp: str


class RequestTypeValidation(Enum):
    IHD = "IHD"
    ANONYMIZED = "Anonymized"
    STUDY = "Keycoded"


# ------------------------------
# 异常处理装饰器核心实现
# ------------------------------
def exception_handler_decorator(func):
    @functools.wraps(func)
    async def wrapper(eventHubMessage: str):
        try:
            # 调用原业务函数,直接传递原始消息参数
            await func(eventHubMessage)
        except json.JSONDecodeError as e:
            # 处理JSON格式错误
            error_payload = {
                "original_eventhub_message": eventHubMessage,
                "error_category": "INVALID_JSON",
                "error_message": str(e),
                "error_type": type(e).__name__,
                "occurred_at": datetime.utcnow().isoformat()
            }
            logging.error(f"JSON解析失败: {json.dumps(error_payload)}")
            await send_error_event(error_payload)
        except ValueError as e:
            # 处理数据校验错误
            error_payload = {
                "original_eventhub_message": eventHubMessage,
                "error_category": "VALIDATION_FAILED",
                "error_message": str(e),
                "error_type": type(e).__name__,
                "occurred_at": datetime.utcnow().isoformat()
            }
            logging.error(f"数据校验失败: {json.dumps(error_payload)}")
            await send_error_event(error_payload)
        except HTTPException as e:
            # 处理下游API调用错误
            error_payload = {
                "original_eventhub_message": eventHubMessage,
                "error_category": "API_CALL_FAILED",
                "error_message": e.detail,
                "status_code": e.status_code,
                "error_type": type(e).__name__,
                "occurred_at": datetime.utcnow().isoformat()
            }
            logging.error(f"接口调用错误: {json.dumps(error_payload)}")
            await send_error_event(error_payload)
        except Exception as e:
            # 兜底处理所有未捕获的异常
            error_payload = {
                "original_eventhub_message": eventHubMessage,
                "error_category": "UNEXPECTED_ERROR",
                "error_message": str(e),
                "error_type": type(e).__name__,
                "stack_trace": traceback.format_exc(),  # 保留堆栈信息,方便排查根因
                "occurred_at": datetime.utcnow().isoformat()
            }
            logging.error(f"未预期的错误: {json.dumps(error_payload)}")
            await send_error_event(error_payload)
    return wrapper


# ------------------------------
# 错误事件发送函数(替换成你的实际逻辑)
# ------------------------------
async def send_error_event(error_payload: dict):
    # 从配置中读取错误处理Event Hub的连接信息
    conn_str = settings_helper.get_value("ERROR_EVENT_HUB_CONNECTION_STRING")
    eventhub_name = settings_helper.get_value("ERROR_EVENT_HUB_NAME")
    
    # 使用Azure Event Hub SDK发送结构化错误事件
    async with AsyncEventHubProducerClient.from_connection_string(
        conn_str, eventhub_name=eventhub_name
    ) as producer:
        event_data_batch = await producer.create_batch()
        event_data_batch.add(EventData(json.dumps(error_payload)))
        await producer.send_batch(event_data_batch)
    logging.info(f"错误事件已发送: {error_payload['error_category']}")


# ------------------------------
# 简化后的业务主函数(移除内部try-except,交给装饰器统一处理)
# ------------------------------
@exception_handler_decorator
async def main(eventHubMessage: str):
    # 反序列化JSON payload
    data = json.loads(eventHubMessage)
    logging.info(f"开始处理消息: {data}")

    if isinstance(data, list):
        # 处理批量消息
        for idx, record in enumerate(data):
            validate_request(record)
            await process_request(record)
            logging.info(f"批量消息第{idx+1}条处理完成")
    else:
        # 处理单条消息
        validate_request(data)
        await process_request(data)

    logging.info("所有消息处理完成")


# 数据校验函数(保留原逻辑,新增枚举值校验)
def validate_request(data: dict):
    required_keys = [
        "RequestID",
        "ProductType",
        "RequestType",
        "RequesterADGroupName",
        "TimeStamp",
    ]
    for key in required_keys:
        if key not in data:
            raise ValueError(f"缺少必填字段: {key}")
    # 新增RequestType的枚举合法性校验
    if data["RequestType"] not in [rt.value for rt in RequestTypeValidation]:
        raise ValueError(f"无效的RequestType: {data['RequestType']},允许值: {[rt.value for rt in RequestTypeValidation]}")


# 业务处理函数(保留原逻辑)
async def process_request(data):
    logging.info(f"处理单个请求: {data}")

    eventhub_request = EventHubRequest(
        IHDRequestID=data["RequestID"],
        ProductType=data["ProductType"],
        RequestType=data["RequestType"],
        RequesterADGroupName=data["RequesterADGroupName"],
        TimeStamp=data["TimeStamp"],
    )

    full_request = eventhub_request.dict()
    api_url = settings_helper.get_value("API_URL") + "access_manager/send_request_v2"
    headers = {
        "content-type": "application/json",
        "x-functions-key": settings_helper.get_value("MASTER_KEY"),
    }

    async with httpx.AsyncClient() as client:
        logging.info(f"调用下游API: {api_url}")
        response = await client.post(
            api_url, json=full_request, headers=headers, timeout=60.0
        )
        response.raise_for_status()
        logging.info(f"API调用成功,状态码: {response.status_code}")

关键优势说明

  1. 原消息100%保留:装饰器直接绑定eventHubMessage参数,异常发生时可以完整打包原始消息内容,不会丢失任何上下文
  2. 异常分类清晰:针对不同类型的异常做了分类处理,错误事件包含明确的错误类别,方便后续的监控和排查
  3. 业务与异常解耦:主函数只需要关注业务逻辑,所有异常处理逻辑都集中在装饰器中,符合单一职责原则
  4. 可扩展性强:如果需要新增异常类型,只需在装饰器中添加对应的except分支;如果要修改错误事件的发送目标,只需调整send_error_event函数

额外优化建议

  • 如果你的批量消息需要失败后继续处理其他记录,可以在main函数的循环中添加局部try-except,并在其中调用send_error_event传递单条失败记录的详情
  • 可以为send_error_event添加重试逻辑(比如用tenacity库),避免临时网络错误导致错误事件丢失
  • 可以在错误事件中加入更多上下文信息,比如当前Azure函数的实例ID、版本号等,方便生产环境调试

备注:内容来源于stack exchange,提问作者BelgratoSystem

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.15 03:24:31