Skip to content

Commit 2f69b41

Browse files
amin-farjadiAmin Farjadileandrodamascena
authored
feat[kafka]: avro offline schema handling (#8442)
* feat: add optional args to 'SchemaConfig' for ignoring prefix * chore: add documentation * address comment on issue * fix(kafka): validate Confluent Avro wire format --------- Co-authored-by: Amin Farjadi <amin.farjadi@eonnext.com> Co-authored-by: Leandro Damascena <lcdama@amazon.pt>
1 parent fe7d5d2 commit 2f69b41

6 files changed

Lines changed: 318 additions & 8 deletions

File tree

‎aws_lambda_powertools/utilities/kafka/consumer_records.py‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -41,17 +41,20 @@ def key(self) -> Any:
4141
schema_type = None
4242
schema_value = None
4343
output_serializer = None
44+
key_schema_wire_format = None
4445

4546
if self.schema_config and self.schema_config.key_schema_type:
4647
schema_type = self.schema_config.key_schema_type
4748
schema_value = self.schema_config.key_schema
4849
output_serializer = self.schema_config.key_output_serializer
50+
key_schema_wire_format = self.schema_config.key_schema_wire_format
4951

5052
# Always use get_deserializer if None it will default to DEFAULT
5153
deserializer = get_deserializer(
5254
schema_type=schema_type,
5355
schema_value=schema_value,
5456
field_metadata=self.key_schema_metadata,
57+
wire_format=key_schema_wire_format,
5558
)
5659
deserialized_value = deserializer.deserialize(key)
5760

@@ -69,19 +72,22 @@ def value(self) -> Any:
6972
schema_type = None
7073
schema_value = None
7174
output_serializer = None
75+
value_schema_wire_format = None
7276

7377
logger.debug("Deserializing value field")
7478

7579
if self.schema_config and self.schema_config.value_schema_type:
7680
schema_type = self.schema_config.value_schema_type
7781
schema_value = self.schema_config.value_schema
7882
output_serializer = self.schema_config.value_output_serializer
83+
value_schema_wire_format = self.schema_config.value_schema_wire_format
7984

8085
# Always use get_deserializer if None it will default to DEFAULT
8186
deserializer = get_deserializer(
8287
schema_type=schema_type,
8388
schema_value=schema_value,
8489
field_metadata=self.value_schema_metadata,
90+
wire_format=value_schema_wire_format,
8591
)
8692
deserialized_value = deserializer.deserialize(value)
8793

‎aws_lambda_powertools/utilities/kafka/deserializer/avro.py‎

Lines changed: 35 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,7 @@
22

33
import io
44
import logging
5-
from typing import Any
5+
from typing import Any, Literal
66

77
from avro.io import BinaryDecoder, DatumReader
88
from avro.schema import parse as parse_schema
@@ -16,6 +16,9 @@
1616

1717
logger = logging.getLogger(__name__)
1818

19+
_CONFLUENT_HEADER_SIZE = 5
20+
_CONFLUENT_MAGIC_BYTE = 0x00
21+
1922

2023
class AvroDeserializer(DeserializerBase):
2124
"""
@@ -25,16 +28,43 @@ class AvroDeserializer(DeserializerBase):
2528
a provided Avro schema definition.
2629
"""
2730

28-
def __init__(self, schema_str: str, field_metadata: dict[str, Any] | None = None):
31+
def __init__(
32+
self,
33+
schema_str: str,
34+
field_metadata: dict[str, Any] | None = None,
35+
wire_format: Literal["CONFLUENT"] | None = None,
36+
):
2937
try:
3038
self.parsed_schema = parse_schema(schema_str)
3139
self.reader = DatumReader(self.parsed_schema)
3240
self.field_metatada = field_metadata
41+
self.wire_format = wire_format
3342
except Exception as e:
3443
raise KafkaConsumerAvroSchemaParserError(
3544
f"Invalid Avro schema. Please ensure the provided avro schema is valid: {type(e).__name__}: {str(e)}",
3645
) from e
3746

47+
def _strip_wire_format_header(self, value: bytes) -> bytes:
48+
if self.wire_format is None:
49+
return value
50+
51+
if self.wire_format != "CONFLUENT":
52+
raise KafkaConsumerDeserializationError(f"Unsupported Avro wire format: {self.wire_format}")
53+
54+
if len(value) < _CONFLUENT_HEADER_SIZE:
55+
raise KafkaConsumerDeserializationError(
56+
"Invalid Confluent wire format: payload must contain a 5-byte header",
57+
)
58+
59+
if value[0] != _CONFLUENT_MAGIC_BYTE:
60+
raise KafkaConsumerDeserializationError(
61+
"Invalid Confluent wire format: expected magic byte 0x00",
62+
)
63+
64+
schema_id = int.from_bytes(value[1:_CONFLUENT_HEADER_SIZE], byteorder="big")
65+
logger.debug("Deserializing Confluent payload with schema ID %s", schema_id)
66+
return value[_CONFLUENT_HEADER_SIZE:]
67+
3868
def deserialize(self, data: bytes | str) -> object:
3969
"""
4070
Deserialize Avro binary data to a Python dictionary.
@@ -75,9 +105,12 @@ def deserialize(self, data: bytes | str) -> object:
75105

76106
try:
77107
value = self._decode_input(data)
108+
value = self._strip_wire_format_header(value)
78109
bytes_reader = io.BytesIO(value)
79110
decoder = BinaryDecoder(bytes_reader)
80111
return self.reader.read(decoder)
112+
except KafkaConsumerDeserializationError:
113+
raise
81114
except Exception as e:
82115
raise KafkaConsumerDeserializationError(
83116
f"Error trying to deserialize avro data - {type(e).__name__}: {str(e)}",

‎aws_lambda_powertools/utilities/kafka/deserializer/deserializer.py‎

Lines changed: 20 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,7 @@
11
from __future__ import annotations
22

33
import hashlib
4-
from typing import TYPE_CHECKING, Any
4+
from typing import TYPE_CHECKING, Any, Literal
55

66
from aws_lambda_powertools.utilities.kafka.deserializer.default import DefaultDeserializer
77
from aws_lambda_powertools.utilities.kafka.deserializer.json import JsonDeserializer
@@ -13,7 +13,12 @@
1313
_deserializer_cache: dict[str, DeserializerBase] = {}
1414

1515

16-
def _get_cache_key(schema_type: str | object, schema_value: Any, field_metadata: dict[str, Any]) -> str:
16+
def _get_cache_key(
17+
schema_type: str | object,
18+
schema_value: Any,
19+
field_metadata: dict[str, Any],
20+
wire_format: Literal["CONFLUENT"] | None,
21+
) -> str:
1722
schema_metadata = None
1823

1924
if field_metadata:
@@ -30,10 +35,15 @@ def _get_cache_key(schema_type: str | object, schema_value: Any, field_metadata:
3035
# For objects like Protobuf, use the object id
3136
schema_hash = f"{str(id(schema_value))}_{schema_metadata}"
3237

33-
return f"{schema_type}_{schema_hash}"
38+
return f"{schema_type}_{schema_hash}_{wire_format}"
3439

3540

36-
def get_deserializer(schema_type: str | object, schema_value: Any, field_metadata: Any) -> DeserializerBase:
41+
def get_deserializer(
42+
schema_type: str | object,
43+
schema_value: Any,
44+
field_metadata: Any,
45+
wire_format: Literal["CONFLUENT"] | None = None,
46+
) -> DeserializerBase:
3747
"""
3848
Factory function to get the appropriate deserializer based on schema type.
3949
@@ -81,7 +91,7 @@ def get_deserializer(schema_type: str | object, schema_value: Any, field_metadat
8191
"""
8292

8393
# Generate a cache key based on schema type and value
84-
cache_key = _get_cache_key(schema_type, schema_value, field_metadata)
94+
cache_key = _get_cache_key(schema_type, schema_value, field_metadata, wire_format)
8595

8696
# Check if we already have this deserializer in cache
8797
if cache_key in _deserializer_cache:
@@ -93,7 +103,11 @@ def get_deserializer(schema_type: str | object, schema_value: Any, field_metadat
93103
# Import here to avoid dependency if not used
94104
from aws_lambda_powertools.utilities.kafka.deserializer.avro import AvroDeserializer
95105

96-
deserializer = AvroDeserializer(schema_str=schema_value, field_metadata=field_metadata)
106+
deserializer = AvroDeserializer(
107+
schema_str=schema_value,
108+
field_metadata=field_metadata,
109+
wire_format=wire_format,
110+
)
97111
elif schema_type == "PROTOBUF":
98112
# Import here to avoid dependency if not used
99113
from aws_lambda_powertools.utilities.kafka.deserializer.protobuf import ProtobufDeserializer

‎aws_lambda_powertools/utilities/kafka/schema_config.py‎

Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -26,6 +26,14 @@ class SchemaConfig:
2626
Schema definition for message keys. Required when key_schema_type is 'AVRO' or 'PROTOBUF'.
2727
key_output_serializer : Any, optional
2828
Custom serializer for message keys. Supports Pydantic classes, Dataclasses and Custom Class
29+
value_schema_wire_format : {'CONFLUENT', None}, default=None
30+
Set this when a Confluent schema-registry-aware serializer produced the value payload
31+
but you are supplying the Avro schema offline rather than using the ESM Schema Registry integration.
32+
Only applies to AVRO values.
33+
key_schema_wire_format : {'CONFLUENT', None}, default=None
34+
Set this when a Confluent schema-registry-aware serializer produced the key payload
35+
but you are supplying the Avro schema offline rather than using the ESM Schema Registry integration.
36+
Only applies to AVRO keys.
2937
3038
Raises
3139
------
@@ -63,21 +71,39 @@ def __init__(
6371
key_schema_type: Literal["AVRO", "PROTOBUF", "JSON"] | None = None,
6472
key_schema: str | None = None,
6573
key_output_serializer: Any | None = None,
74+
value_schema_wire_format: Literal["CONFLUENT"] | None = None,
75+
key_schema_wire_format: Literal["CONFLUENT"] | None = None,
6676
):
6777
# Validate schema requirements
6878
self._validate_schema_requirements(value_schema_type, value_schema, "value")
6979
self._validate_schema_requirements(key_schema_type, key_schema, "key")
80+
self._validate_wire_format(value_schema_wire_format, value_schema_type, "value")
81+
self._validate_wire_format(key_schema_wire_format, key_schema_type, "key")
7082

7183
self.value_schema_type = value_schema_type
7284
self.value_schema = value_schema
7385
self.value_output_serializer = value_output_serializer
7486
self.key_schema_type = key_schema_type
7587
self.key_schema = key_schema
7688
self.key_output_serializer = key_output_serializer
89+
self.value_schema_wire_format = value_schema_wire_format
90+
self.key_schema_wire_format = key_schema_wire_format
7791

7892
def _validate_schema_requirements(self, schema_type: str | None, schema: str | None, prefix: str) -> None:
7993
"""Validate that schema is provided when required by schema_type."""
8094
if schema_type in ["AVRO", "PROTOBUF"] and schema is None:
8195
raise KafkaConsumerMissingSchemaError(
8296
f"{prefix}_schema must be provided when {prefix}_schema_type is {schema_type}",
8397
)
98+
99+
def _validate_wire_format(self, wire_format: str | None, schema_type: str | None, prefix: str) -> None:
100+
"""Validate the wire format for a key or value payload."""
101+
102+
if wire_format is None:
103+
return
104+
105+
if wire_format != "CONFLUENT":
106+
raise ValueError(f"{prefix}_schema_wire_format must be 'CONFLUENT'.")
107+
108+
if schema_type != "AVRO":
109+
raise ValueError(f"{prefix}_schema_wire_format is supported only when {prefix}_schema_type is 'AVRO'.")

‎docs/utilities/kafka.md‎

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,7 @@ flowchart LR
2929
* Support for key and value deserialization
3030
* Support for custom output serializers (e.g., dataclasses, Pydantic models)
3131
* Support for ESM with and without Schema Registry integration
32+
* Support for offline Avro schemas with schema-registry wire-format prefixes (Confluent only)
3233
* Proper error handling for deserialization issues
3334

3435
## Terminology
@@ -255,6 +256,45 @@ Each Kafka record contains important metadata that you can access alongside the
255256
| `value_schema_metadata` | Metadata about the value schema like `schemaId` and `dataFormat` | Data format and schemaId propagated when integrating with Schema Registry |
256257
| `key_schema_metadata` | Metadata about the key schema like `schemaId` and `dataFormat` | Data format and schemaId propagated when integrating with Schema Registry |
257258

259+
### Using an offline Avro schema with a schema-registry wire-format prefix
260+
261+
When Confluent serializes messages with its schema-registry-aware Avro serializer (for example, `KafkaAvroSerializer`), each payload carries a wire-format header before the Avro body.
262+
The header is 5 bytes long: 1-byte magic byte (`0x00`) followed by a 4-byte big-endian schema ID.
263+
264+
When the ESM Schema Registry integration is enabled, Lambda strips those bytes and populates the record's schema metadata. When you use an **offline Avro schema** without the ESM Schema Registry integration, the header reaches the function and prevents plain Avro deserialization.
265+
266+
Set `value_schema_wire_format` or `key_schema_wire_format` on `SchemaConfig` to `"CONFLUENT"`. Powertools validates the magic byte and strips the 5-byte header before running the Avro decoder.
267+
268+
???+ info "When do I need this?"
269+
Use this option when you supply the Avro schema and the producer uses the Confluent wire format. If ESM Schema Registry integration has already removed the header, leave the option as `None`.
270+
271+
=== "Offline Avro schema with a Confluent prefix"
272+
273+
```python hl_lines="10"
274+
from aws_lambda_powertools.utilities.kafka import SchemaConfig, kafka_consumer
275+
from aws_lambda_powertools.utilities.kafka.consumer_records import ConsumerRecords
276+
from aws_lambda_powertools.utilities.typing import LambdaContext
277+
278+
AVRO_SCHEMA = open("user.avsc").read()
279+
280+
schema_config = SchemaConfig(
281+
value_schema_type="AVRO",
282+
value_schema=AVRO_SCHEMA,
283+
value_schema_wire_format="CONFLUENT",
284+
)
285+
286+
287+
@kafka_consumer(schema_config=schema_config)
288+
def lambda_handler(event: ConsumerRecords, context: LambdaContext):
289+
for record in event.records:
290+
# record.value is the deserialized Avro payload
291+
# with the validated 5-byte wire-format header removed.
292+
...
293+
```
294+
295+
???+ warning "Scope"
296+
`value_schema_wire_format` and `key_schema_wire_format` apply only to **Avro** payloads. Leave them as `None` when ESM Schema Registry integration has already removed the wire-format header.
297+
258298
### Custom output serializers
259299

260300
Transform deserialized data into your preferred object types using output serializers. This can help you integrate Kafka data with your domain models and application architecture, providing type hints, validation, and structured data access.

0 commit comments

Comments
 (0)