为何Kafka触发的Nuclio函数中事件偏移量始终为0?
Kafka触发器中event对象offset始终为0的原因分析
我编写了一个简单的handler函数来探究Kafka生成的event对象,部署后发送了若干Kafka消息。
def handler(context, event): try: print(f"Variables inside the event object: {vars(event)}") print(f"Attributes of event object: {dir(event)}") except Exception as e: print(f"Exception {e} occured")
在Pod日志中,每次Kafka事件被触发时,输出如下:
Variables inside the event object: {'body': b'{"value0":"value","value1":[1,2,3]}', 'content_type': b'', 'trigger': <nuclio_sdk.event.TriggerInfo object at 0x7f928c5c3040>, 'fields': {}, 'headers': {}, 'id': b'd6db95ca-0659-490b-b8f9-4f7dc4b09b70', 'method': b'', 'path': b'event_object_topic', 'size': 35, 'timestamp': datetime.datetime(2024, 3, 8, 8, 0, 14), 'url': b'', 'shard_id': 0, 'num_shards': 0, 'type': b'', 'type_version': b'', 'version': b'', 'last_in_batch': False, 'offset': 0} Attributes of event object: ['__class__', '__delattr__', '__dict__', '__dir__', '__doc__', '__eq__', '__format__', '__ge__', '__getattribute__', '__gt__', '__hash__', '__init__', '__init_subclass__', '__le__', '__lt__', '__module__', '__ne__', '__new__', '__reduce__', '__reduce_ex__', '__repr__', '__setattr__', '__sizeof__', '__str__', '__subclasshook__', '__weakref__', 'body', 'content_type', 'deserialize', 'fields', 'from_json', 'from_msgpack', 'get_header', 'headers', 'id', 'last_in_batch', 'method', 'num_shards', 'offset', 'path', 'shard_id', 'size', 'timestamp', 'to_json', 'trigger', 'type', 'type_version', 'url', 'version'] {"level":"debug","time":"2024-03-08T08:00:14.940Z","name":"processor.kafka-cluster.epic-others-event-loader.sarama","message":"Sarama: client/metadata fetching metadata for [event_object_topic] from broker vmwxvsapp05-xxxxxxxxxxxx:9092"} Variables inside the event object: {'body': b'{"value0":"value","value1":[1,2,3]}', 'content_type': b'', 'trigger': <nuclio_sdk.event.TriggerInfo object at 0x7f6169cfe040>, 'fields': {}, 'headers': {}, 'id': b'cd8cde08-b70f-4f34-abee-c936c3e6acf9', 'method': b'', 'path': b'event_object_topic', 'size': 35, 'timestamp': datetime.datetime(2024, 3, 8, 8, 0, 15), 'url': b'', 'shard_id': 0, 'num_shards': 0, 'type': b'', 'type_version': b'', 'version': b'', 'last_in_batch': False, 'offset': 0} Attributes of event object: ['__class__', '__delattr__', '__dict__', '__dir__', '__doc__', '__eq__', '__format__', '__ge__', '__getattribute__', '__gt__', '__hash__', '__init__', '__init_subclass__', '__le__', '__lt__', '__module__', '__ne__', '__new__', '__reduce__', '__reduce_ex__', '__repr__', '__setattr__', '__sizeof__', '__str__', '__subclasshook__', '__weakref__', 'body', 'content_type', 'deserialize', 'fields', 'from_json', 'from_msgpack', 'get_header', 'headers', 'id', 'last_in_batch', 'method', 'num_shards', 'offset', 'path', 'shard_id', 'size', 'timestamp', 'to_json', 'trigger', 'type', 'type_version', 'url', 'version'] {"level":"debug","time":"2024-03-08T08:00:15.938Z","name":"processor.kafka-cluster.epic-others-event-loader.sarama","message":"Sarama: client/metadata fetching metadata for [event_object_topic] from broker vmwxvsapp05-xxxxxxxxxxxx:9092"} Variables inside the event object: {'body': b'{"value0":"value","value1":[1,2,3]}', 'content_type': b'', 'trigger': <nuclio_sdk.event.TriggerInfo object at 0x7f928c5c3460>, 'fields': {}, 'headers': {}, 'id': b'e83bd16c-1d3c-490a-b3f1-187ae99617f8', 'method': b'', 'path': b'event_object_topic', 'size': 35, 'timestamp': datetime.datetime(2024, 3, 8, 8, 0, 18), 'url': b'', 'shard_id': 0, 'num_shards': 0, 'type': b'', 'type_version': b'', 'version': b'', 'last_in_batch': False, 'offset': 0} Attributes of event object: ['__class__', '__delattr__', '__dict__', '__dir__', '__doc__', '__eq__', '__format__', '__ge__', '__getattribute__', '__gt__', '__hash__', '__init__', '__init_subclass__', '__le__', '__lt__', '__module__', '__ne__', '__new__', '__reduce__', '__reduce_ex__', '__repr__', '__setattr__', '__sizeof__', '__str__', '__subclasshook__', '__weakref__', 'body', 'content_type', 'deserialize', 'fields', 'from_json', 'from_msgpack', 'get_header', 'headers', 'id', 'last_in_batch', 'method', 'num_shards', 'offset', 'path', 'shard_id', 'size', 'timestamp', 'to_json', 'trigger', 'type', 'type_version', 'url', 'version']
从输出可见,offset始终为0,但在Kafka Control Center及对应主题中,每条消息的偏移量均不相同,这是什么原因?
原因及解决方案
这是因为Nuclio的Kafka触发器默认没有将Kafka消息的原生offset映射到event对象的offset字段,这个字段是为兼容其他触发器场景设计的,并非专门承载Kafka偏移量。
Kafka的消息元数据(包括offset、分区等)都存储在event.trigger这个TriggerInfo对象里,需要从这里提取才能拿到正确的偏移量。
修改后的handler函数示例:
def handler(context, event): try: print(f"Variables inside the event object: {vars(event)}") print(f"Attributes of event object: {dir(event)}") # 获取Kafka原生消息偏移量 if hasattr(event.trigger, 'offset'): print(f"Kafka message offset: {event.trigger.offset}") # 获取Kafka消息所在分区 if hasattr(event.trigger, 'partition'): print(f"Kafka partition: {event.trigger.partition}") except Exception as e: print(f"Exception {e} occurred")
部署修改后的代码后,日志中就能看到和Kafka Control Center一致的偏移量了。另外日志中num_shards显示为0也是同样道理,这个字段对应Nuclio触发器的分片数,而非Kafka的分区数,Kafka分区数同样需要从event.trigger中获取。
内容的提问来源于stack exchange,提问作者fahadhub
相关产品推荐
相关产品推荐

