아래 코드는 브로커로부터 이벤트를 구독해, 해당 이벤트가 발생했을 때 다음 단계의 작업(예 리스크 분석)을 시작합니다. Consumer 에이전트 역시 특정 에이전트에게 호출되지 않고, 이벤트를 보고 스스로 필요 여부를 판단해 행동한다는 점에서 A2A의 자율성이 드러납니다.
# consumer.py (구독 에이전트)
from kafka import KafkaConsumer
import json
consumer = KafkaConsumer(
"a2a-events",
bootstrap_servers=['localhost:9092'],
value_deserializer=lambda v: json.loads(v.decode('utf-8')),
group_id="risk-analyzer",
auto_offset_reset="earliest"
)
for message in consumer:
event = message.value
if event["event_type"] == "analysis.completed":
print(f"[리스크 분석 에이전트] 문서:{event['document_id']} 분석 시작합니다.")
# 추가 분석 로직 수행