[Feature] Support raw Kafka record (KafkaRecord) binding with Protobuf deserialization
评估
调研方向
先从 azure-functions-java-library 中 KafkaRecord.java、KafkaHeader.java、KafkaTimestamp.java 和 KafkaTimestampType.java 的新增内容开始,然后检查 azure-functions-java-worker 中的 KafkaRecordProto.proto、RpcModelBindingDataSource.java 和 KafkaRecordProtoDeserializer.java。验证 protobuf 生成以及 application/json 和 application/x-protobuf 两者的 dispatch,然后测试现有 bindings 是否继续正常工作,并确认 raw records 是否暴露指定的元数据。
由索引模型根据 Issue 内容生成。
描述
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. The host-side Kafka Extension 4.3.1 and .NET Isolated Worker (PR #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-encodedKafkaRecordProto
Currently, Java users can only bind to String, byte[], or Map<String, String> — they cannot access structured metadata like headers, timestamps, or partition info.
Protobuf Schema (shared across all languages)
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
// 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
- POJO types in
java-library(no dependency on host release) - Protobuf deserializer in
java-worker(requires host extension 4.3.1 NuGet — already released)
Related Issues
- Parent: Azure/azure-functions-kafka-extension#612
- .NET Worker (done): Azure/azure-functions-dotnet-worker#3356
- 主要语言
- Java
- 星标
- 103
- 派生
- 74
- 平均合并
- 1 天 11 小时
- 30 天内合并 PR
- 3
环境准备
- 没有 Dockerfile 或 Docker Compose 文件
- 有 Pull Request 模板
- 阅读贡献指南
从这里开始
- 先读完整个 Issue,再读项目的贡献指南。
- 在 Issue 下留言说明你要接手 —— 这能避免两个人做同样的事。
- Fork 仓库,在一个分支上完成修改。
- 提交 Pull Request,并在描述里引用这个 Issue 编号。
Azure/azure-functions-java-worker 的其他 Issue
-
[Feature] Support for reverse propagation of request telemetry tags可能已有人在做 @arnabnandy7 于 60 天前认领。 未关闭
难度 4/5 3-5 天 新手友好度 48/100
-
[Feature] Support OpenTelemetry Baggage可能已有人在做 @arnabnandy7 于 60 天前认领。 未关闭
难度 4/5 3-5 天 新手友好度 48/100
-
mcp-extension
难度 5/5 一周以上 新手友好度 25/100
-
Add advanced MCP tool samples可能重新可做 @ahmedmuhsin 于 176 天前认领,目前没有进行中的 PR。 未关闭mcp-extension
Azure/azure-functions-java-worker#863 · 已指派 1 人 ·
-
Add MCP (Model Context Protocol) resource samples可能重新可做 @ahmedmuhsin 于 176 天前认领,目前没有进行中的 PR。 未关闭mcp-extension
Azure/azure-functions-java-worker#862 · 已指派 1 人 ·
查看 Azure/azure-functions-java-worker 的全部 Issue
相似的 Issue
-
难度 2/5 1-3 小时 新手友好度 88/100
java-native-access/jna#1740 ·
-
难度 1/5 1 小时以内 新手友好度 82/100
-
难度 1/5 1 小时以内 新手友好度 88/100
portfolio-performance/portfolio#6119 ·
维护者通常 8 天内回复
-
难度 2/5 1-3 小时 新手友好度 68/100
jenkinsci/ec2-plugin#2041 ·
-
L: github:actions L: php:composer
难度 2/5 1-3 小时 新手友好度 88/100
dependabot/dependabot-core#16493 ·
维护者通常 1 天内回复