From dbb60383e22168ef1a6cf34485be6029b82cb319 Mon Sep 17 00:00:00 2001 From: Sanjeev Kulkarni Date: Wed, 8 May 2019 16:30:39 -0700 Subject: [PATCH 1/3] Expose MessageId as part of Record interface --- .../main/java/org/apache/pulsar/functions/api/Record.java | 7 +++++++ .../org/apache/pulsar/functions/source/PulsarRecord.java | 3 +++ 2 files changed, 10 insertions(+) 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..cbaa23ca8956d 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 getId() { + 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..e1c128fbfa231 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 getId() { return Optional.of(message.getMessageId().toByteArray()); } } From 3711eef43473675034adc4cef25995fc445f9d1c Mon Sep 17 00:00:00 2001 From: Sanjeev Kulkarni Date: Wed, 8 May 2019 21:32:06 -0700 Subject: [PATCH 2/3] Fixed build --- .../apache/pulsar/io/canal/CanalAbstractSource.java | 12 ++++++------ 1 file changed, 6 insertions(+), 6 deletions(-) diff --git a/pulsar-io/canal/src/main/java/org/apache/pulsar/io/canal/CanalAbstractSource.java b/pulsar-io/canal/src/main/java/org/apache/pulsar/io/canal/CanalAbstractSource.java index c1bcb3d504cb2..c6b75d0bd3d94 100644 --- a/pulsar-io/canal/src/main/java/org/apache/pulsar/io/canal/CanalAbstractSource.java +++ b/pulsar-io/canal/src/main/java/org/apache/pulsar/io/canal/CanalAbstractSource.java @@ -135,7 +135,7 @@ protected void process() { } else { if (flatMessages != null) { CanalRecord canalRecord = new CanalRecord<>(connector); - canalRecord.setId(batchId); + canalRecord.setRecordId(batchId); canalRecord.setRecord(extractValue(flatMessages)); consume(canalRecord); } @@ -159,7 +159,7 @@ protected void process() { static private class CanalRecord implements Record { private V record; - private Long id; + private Long recordId; private CanalConnector connector; public CanalRecord(CanalConnector connector) { @@ -168,7 +168,7 @@ public CanalRecord(CanalConnector connector) { @Override public Optional getKey() { - return Optional.of(Long.toString(id)); + return Optional.of(Long.toString(recordId)); } @Override @@ -178,13 +178,13 @@ public V getValue() { @Override public Optional getRecordSequence() { - return Optional.of(id); + return Optional.of(recordId); } @Override public void ack() { - log.info("CanalRecord ack id is {}", this.id); - connector.ack(this.id); + log.info("CanalRecord ack id is {}", this.recordId); + connector.ack(this.recordId); } } From c132f8b01221759d24ce003b7ea5171c458aaf98 Mon Sep 17 00:00:00 2001 From: Sanjeev Kulkarni Date: Thu, 9 May 2019 04:27:43 -0700 Subject: [PATCH 3/3] Fixed build --- .../java/org/apache/pulsar/functions/api/Record.java | 2 +- .../apache/pulsar/functions/source/PulsarRecord.java | 2 +- .../apache/pulsar/io/canal/CanalAbstractSource.java | 12 ++++++------ 3 files changed, 8 insertions(+), 8 deletions(-) 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 cbaa23ca8956d..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 @@ -44,7 +44,7 @@ default Optional getKey() { /** * Return a record id if the record has one associated. */ - default Optional getId() { + default Optional getRecordId() { return Optional.empty(); } 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 e1c128fbfa231..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 @@ -108,5 +108,5 @@ public void fail() { } @Override - public Optional getId() { return Optional.of(message.getMessageId().toByteArray()); } + public Optional getRecordId() { return Optional.of(message.getMessageId().toByteArray()); } } diff --git a/pulsar-io/canal/src/main/java/org/apache/pulsar/io/canal/CanalAbstractSource.java b/pulsar-io/canal/src/main/java/org/apache/pulsar/io/canal/CanalAbstractSource.java index c6b75d0bd3d94..c1bcb3d504cb2 100644 --- a/pulsar-io/canal/src/main/java/org/apache/pulsar/io/canal/CanalAbstractSource.java +++ b/pulsar-io/canal/src/main/java/org/apache/pulsar/io/canal/CanalAbstractSource.java @@ -135,7 +135,7 @@ protected void process() { } else { if (flatMessages != null) { CanalRecord canalRecord = new CanalRecord<>(connector); - canalRecord.setRecordId(batchId); + canalRecord.setId(batchId); canalRecord.setRecord(extractValue(flatMessages)); consume(canalRecord); } @@ -159,7 +159,7 @@ protected void process() { static private class CanalRecord implements Record { private V record; - private Long recordId; + private Long id; private CanalConnector connector; public CanalRecord(CanalConnector connector) { @@ -168,7 +168,7 @@ public CanalRecord(CanalConnector connector) { @Override public Optional getKey() { - return Optional.of(Long.toString(recordId)); + return Optional.of(Long.toString(id)); } @Override @@ -178,13 +178,13 @@ public V getValue() { @Override public Optional getRecordSequence() { - return Optional.of(recordId); + return Optional.of(id); } @Override public void ack() { - log.info("CanalRecord ack id is {}", this.recordId); - connector.ack(this.recordId); + log.info("CanalRecord ack id is {}", this.id); + connector.ack(this.id); } }