Azure / Azure/azure-functions-java-worker
[Feature] Support raw Kafka record (KafkaRecord) binding with Protobuf deserialization
- Ngôn ngữ chính
- Java
- Star
- 103
- Fork
- 74
- Merge trung bình
- 4 ngày 8 giờ
- Pull request đã merge (30 ngày)
- 2
Mô tả
## Summary
Add support for binding to raw Apache Kafka records (`KafkaRecord` type) in the Java worker, enabling users to access full Kafka message metadata (topic, partition, offset, key/value as raw bytes, headers, timestamp, leader epoch).
This is the Java implementation of [Azure/azure-functions-kafka-extension#612](https://github.com/Azure/azure-functions-kafka-extension/issues/612). The host-side Kafka Extension 4.3.1 and .NET Isolated Worker ([PR #3356](https://github.com/Azure/azure-functions-dotnet-worker/pull/3356)) are already complete.
## Background
The host-side Kafka Extension (4.3.1) serializes `IKafkaEventData` to Protobuf and sends it as `ParameterBindingData` with:
- `source`: `"AzureKafkaRecord"`
- `content_type`: `"application/x-protobuf"`
- `content`: Protobuf-encoded `KafkaRecordProto`
Currently, Java users can only bind to `String`, `byte[]`, or `Map` — they cannot access structured metadata like headers, timestamps, or partition info.
## Protobuf Schema (shared across all languages)
```protobuf
message KafkaRecordProto {
string topic = 1;
int32 partition = 2;
int64 offset = 3;
optional bytes key = 4;
optional bytes value = 5;
KafkaTimestampProto timestamp = 6;
repeated KafkaHeaderProto headers = 7;
optional int32 leader_epoch = 8;
reserved 9 to 15;
}
message KafkaTimestampProto {
int64 unix_timestamp_ms = 1;
int32 type = 2; // 0=NotAvailable, 1=CreateTime, 2=LogAppendTime
}
message KafkaHeaderProto {
string key = 1;
optional bytes value = 2;
}
```
## Required Changes
### Part A: New POJO types (in `azure-functions-java-library`)
| File | Description |
|------|-------------|
| `KafkaRecord.java` | Main POJO: topic, partition, offset, key (byte[]), value (byte[]), timestamp, headers, leaderEpoch (Integer) |
| `KafkaHeader.java` | Header: key (String) + value (byte[]) + `getValueAsString()` helper |
| `KafkaTimestamp.java` | Timestamp: unixTimestampMs (long) + type (enum) + `getDateTime()` -> OffsetDateTime |
| `KafkaTimestampType.java` | Enum: NotAvailable(0), CreateTime(1), LogAppendTime(2) |
No annotation changes needed — existing `@KafkaTrigger` works as-is.
### Part B: Worker-side Protobuf deserializer (in `azure-functions-java-worker`)
| File | Change |
|------|--------|
| `KafkaRecordProto.proto` | New: Add proto schema, configure `protobuf-maven-plugin` for code generation |
| `RpcModelBindingDataSource.java` | Modify: Add `content_type` dispatch — `application/json` -> existing JSON path, `application/x-protobuf` -> new Protobuf deserialization |
| `KafkaRecordProtoDeserializer.java` | New: Map `KafkaRecordProto` -> `KafkaRecord` POJO |
Note: `protobuf-java` 3.25.5 is already in pom.xml. Maven protobuf plugin setup is needed for `.proto` compilation.
### Part C: No changes to `azure-functions-java-additions`
`KafkaRecord` is a data container, not an Azure SDK client — the SdkType/Hydrator pattern is not applicable.
## User Experience
```java
// Existing (continues to work)
@FunctionName("ExistingTrigger")
public void run(@KafkaTrigger(...) String message) { }
// NEW: Full record access
@FunctionName("KafkaRecordTrigger")
public void run(
@KafkaTrigger(name = "record", topic = "my-topic",
brokerList = "%BrokerList%", consumerGroup = "$Default")
KafkaRecord record,
final ExecutionContext context) {
context.getLogger().info("Topic: " + record.getTopic());
context.getLogger().info("Partition: " + record.getPartition());
context.getLogger().info("Offset: " + record.getOffset());
context.getLogger().info("Key: " + new String(record.getKey()));
context.getLogger().info("Timestamp: " + record.getTimestamp().getDateTime());
for (KafkaHeader header : record.getHeaders()) {
context.getLogger().info("Header: " + header.getKey() + " = " + header.getValueAsString());
}
}
// NEW: Batch mode
@FunctionName("KafkaBatchTrigger")
public void run(
@KafkaTrigger(..., cardinality = Cardinality.MANY)
KafkaRecord[] records) { ... }
```
## Breaking Changes
None. This is purely additive. All existing binding types (`String`, `byte[]`, POJO) continue to work.
## Implementation Order
1. POJO types in `java-library` (no dependency on host release)
2. Protobuf deserializer in `java-worker` (requires host extension 4.3.1 NuGet — already released)
## Related Issues
- Parent: [Azure/azure-functions-kafka-extension#612](https://github.com/Azure/azure-functions-kafka-extension/issues/612)
- .NET Worker (done): [Azure/azure-functions-dotnet-worker#3356](https://github.com/Azure/azure-functions-dotnet-worker/pull/3356)
Hướng dẫn đóng góp
Hướng nghiên cứu
Bắt đầu với các phần bổ sung KafkaRecord.java, KafkaHeader.java, KafkaTimestamp.java và KafkaTimestampType.java trong azure-functions-java-library, sau đó kiểm tra KafkaRecordProto.proto, RpcModelBindingDataSource.java và KafkaRecordProtoDeserializer.java trong azure-functions-java-worker. Xác minh việc sinh protobuf và cơ chế dispatch cho cả application/json và application/x-protobuf, sau đó kiểm thử để bảo đảm các binding hiện có vẫn tiếp tục hoạt động và raw records cung cấp các metadata được chỉ định.
Do mô hình lập chỉ mục viết ra từ nội dung của issue.
Đánh giá
- Công nghệ
- java
- Lĩnh vực
- api, backend
- Loại issue
- Tính năng
- Độ khó
- 4/5
- Thời gian dự kiến
- 3-5 ngày
- Mức độ hoạt động
- Ít trao đổi
- Độ rõ ràng
- Đặc tả rõ ràng
- Mức phù hợp với người mới
- 50/100