58 / 89 · 04 API & Contract Testing with AI · Testing Kafka Consumers← prev⊞ allnext →☰ Read as one page
9.4Malformed Event Handling
def test_malformed_event_does_not_crash(self, kafka_producer, order_consumer):
"""Consumer should handle malformed events gracefully."""
# Send invalid JSON
kafka_producer.produce(
'order-events',
key="bad",
value=b"this is not json"
)
kafka_producer.flush()
# Consumer should not crash
order_consumer.poll(timeout=5.0)
assert order_consumer.is_alive(), "Consumer crashed on malformed event"
def test_event_with_missing_required_fields(self, kafka_producer, order_consumer, db):
"""Consumer should reject events missing required fields."""
incomplete_event = {
"event_type": "order.created",
"data": {
"order_id": "ord-incomplete"
# Missing: customer_id, total, items
}
}
kafka_producer.produce(
'order-events',
key="ord-incomplete",
value=json.dumps(incomplete_event).encode('utf-8')
)
kafka_producer.flush()
time.sleep(5)
order_consumer.poll(timeout=1.0)
order = db.get_order("ord-incomplete")
assert order is None, "Incomplete event should not create an order"
def test_event_with_unknown_type(self, kafka_producer, order_consumer):
"""Consumer should skip events with unrecognized event_type."""
unknown_event = {
"event_type": "order.teleported", # Not a real event type
"data": {"order_id": "ord-unknown"}
}
kafka_producer.produce(
'order-events',
key="ord-unknown",
value=json.dumps(unknown_event).encode('utf-8')
)
kafka_producer.flush()
order_consumer.poll(timeout=5.0)
assert order_consumer.is_alive(), "Consumer crashed on unknown event type"