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

尝试向Event Hub批量发送事件时出现AttributeError错误

问题:Azure Function批量发送Event Hub时出现AttributeError错误

运行从Blob拉取多行JSON事件并批量发送到Event Hub的Azure Function时,出现以下错误:

Result: Failure Exception: AttributeError: 'list' object has no attribute 'header' Stack: File "/azure-functions-host/workers/python/3.11/LINUX/X64/azure_functions_worker/dispatcher.py", line 491, in _handle__invocation_request call_result = await self._run_async_func( ^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/azure-functions-host/workers/python/3.11/LINUX/X64/azure_functions_worker/dispatcher.py", line 774, in _run_async_func return await ExtensionManager.get_async_invocation_wrapper( ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/azure-functions-host/workers/python/3.11/LINUX/X64/azure_functions_worker/extension.py", line 147, in get_async_invocation_wrapper result = await function(**args) ^^^^^^^^^^^^^^^^^^^^^^ File "/home/site/wwwroot/AzureFunctionCloudflare/main-edited.py", line 83, in main await asyncio.gather(*cors) File "/home/site/wwwroot/AzureFunctionCloudflare/main-edited.py", line 189, in process_blob await event_data_batch.add([EventData(json.dumps(event))]) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/home/site/wwwroot/.python_packages/lib/site-packages/azure/eventhub/_common.py", line 685, in add outgoing_event_data = transform_outbound_single_message( ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/home/site/wwwroot/.python_packages/lib/site-packages/azure/eventhub/_utils.py", line 261, in transform_outbound_single_message amqp_message = to_outgoing_amqp_message(message) ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "/home/site/wwwroot/.python_packages/lib/site-packages/azure/eventhub/_transport/_pyamqp_transport.py", line 99, in to_outgoing_amqp_message header_vals = annotated_message.header.values() if annotated_message.header else None ^^^^^^^^^^^^^^^^^^^^^^^^

Blob中的数据为多行独立JSON,示例如下:

{"QueryName":"domain1.co.uk.","QueryType":2,"ResponseCode":0,"Timestamp":"2023-10-24T15:28:50Z","SourceIP":"x.x.x.x"}
{"QueryName":"domain2.co.uk.","QueryType":2,"ResponseCode":0,"Timestamp":"2023-10-24T15:28:50Z","SourceIP":"x.x.x.x"}
{"QueryName":"domain3.co.uk.","QueryType":2,"ResponseCode":0,"Timestamp":"2023-10-24T15:28:50Z","SourceIP":"x.x.x.x"}

处理Blob的函数代码片段:

async def process_blob(self, blob, container_client, session: aiohttp.ClientSession):
        
    async with self.semaphore:
        logging.info("Start processing {}".format(blob['name']))
        blob_client = container_client.get_blob_client(blob['name'])
        tags = await blob_client.get_blob_tags()
        updated_tags = {'StartedProcessing': 'true'}
        tags.update(updated_tags)
        await blob_client.set_blob_tags(tags)
    
        # Create a producer client to send messages to the Event Hub
        connection_str = 'removed'
        event_hub_producer = EventHubProducerClient.from_connection_string(connection_str, eventhub_name='cloudflare')
    
        blob_cor = await container_client.download_blob(blob['name'])
        s = ''
    
        async with event_hub_producer:
    
            # Create a Batch
            event_data_batch = await event_hub_producer.create_batch()
    
            # Process blob chunks
            async for chunk in blob_cor.chunks():
                s += chunk.decode()
                lines = re.split(r'{0}'.format(LINE_SEPARATOR), s)
                lines = self.merge_lines(data=lines)
                for n, line in enumerate(lines):
                    if n < len(lines) - 1:
                        if line:
                            try:
                                event = json.loads(line)
                                await event_data_batch.add([EventData(json.dumps(event))])
                            except JSONDecodeError as je:
                                logging.error(
                                    'JSONDecode error while loading json event at line value {}. blob name: {}. Error {}'.format(
                                        line, blob['name'], str(je)))
                            except ValueError as e:
                                logging.error(
                                    'Error while loading json Event at line value {}. blob name: {}. Error: {}'.format(
                                        line, blob['name'], str(e)))
                                raise e
                    s = line
            if s:
                try:
                    event = json.loads(s)
                except JSONDecodeError as je:
                    logging.error(
                        'JSONDecode error while loading json event at line value {}. blob name: {}. Error {}'.format(
                            line, blob['name'], str(je)))
                except ValueError as e:
                    logging.error('Error while loading json Event at s value {}. blob name: {}. Error: {}'.format(
                        line, blob['name'], str(e)))
                    raise e
                await event_data_batch.add([EventData(json.dumps(event))])
    
            # Send the batch of Events
            await event_hub_producer.send_batch(event_data_batch)
    
            # Delete the Blob
            await self.delete_blob(blob, container_client)
            self.total_blobs += 1

问题分析与解决方案

错误根源

报错的核心原因是EventDataBatch.add()方法被传入了列表类型的参数,但该方法设计为接受单个EventData对象。代码中await event_data_batch.add([EventData(json.dumps(event))])把EventData包装在了列表里,导致内部处理逻辑将这个列表当作单个消息对象,试图访问它的header属性,而列表并没有该属性,因此抛出AttributeError。

修改方案

将两处调用add方法的代码中的列表包裹移除,直接传入单个EventData对象:

  1. 循环内的调用修改为:
await event_data_batch.add(EventData(json.dumps(event)))
  1. 最后处理剩余内容的调用修改为:
await event_data_batch.add(EventData(json.dumps(event)))

补充说明

如果需要批量添加多个EventData对象,可以使用EventDataBatch.add_batch()方法,或者通过循环逐个调用add()。在当前场景下,代码原本就是逐行处理单个JSON事件,因此直接传入单个EventData即可满足需求。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 14:22:34