Azure / Azure/azure-functions-java-worker

[Feature] Support raw Kafka record (KafkaRecord) binding with Protobuf deserialization

Open
#868 0 comments 0 reactions 0 assignees View on GitHub
Dominant language
Java
Stars
103
Forks
74
Avg merge
4d 8h
Merged PRs (30d)
2

Description

## 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)

Contributor guide

Open the contributing guide

Assessment

This issue has not been assessed yet.

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.