diff --git a/pulsar-functions/api-java/src/main/java/org/apache/pulsar/functions/api/Record.java b/pulsar-functions/api-java/src/main/java/org/apache/pulsar/functions/api/Record.java index b99c4dc0e746b..74142946aa2f2 100644 --- a/pulsar-functions/api-java/src/main/java/org/apache/pulsar/functions/api/Record.java +++ b/pulsar-functions/api-java/src/main/java/org/apache/pulsar/functions/api/Record.java @@ -41,6 +41,13 @@ default Optional getKey() { return Optional.empty(); } + /** + * Return a record id if the record has one associated. + */ + default Optional getRecordId() { + return Optional.empty(); + } + /** * Retrieves the actual data of the record. * diff --git a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/source/PulsarRecord.java b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/source/PulsarRecord.java index c03ed9fa834a9..96c24c8ec41bf 100644 --- a/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/source/PulsarRecord.java +++ b/pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/source/PulsarRecord.java @@ -106,4 +106,7 @@ public void ack() { public void fail() { this.failFunction.run(); } + + @Override + public Optional getRecordId() { return Optional.of(message.getMessageId().toByteArray()); } }