This is an automated email from the ASF dual-hosted git repository.
oscerd pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/main by this push:
new e306cfd8ceeb CAMEL-24935: camel-pulsar - report the message id as a
header, not as the body (#26783)
e306cfd8ceeb is described below
commit e306cfd8ceeb0dfcbaabcf768551e67093fd1cd7
Author: Andrea Cosentino <[email protected]>
AuthorDate: Fri Sep 25 10:49:38 2026 +0200
CAMEL-24935: camel-pulsar - report the message id as a header, not as the
body (#26783)
* CAMEL-24935: camel-pulsar - report the message id as a header, not as the
body
After a successful send the producer replaced the message body with the
MessageId returned by the
broker, so every step after the to("pulsar:...") saw a MessageId instead of
the payload, and
.to("pulsar:a").to("pulsar:b") published a Java-serialized MessageId to the
second topic, since
serialize() falls back to Java serialization when no type converter
applies. Nothing documented it and
no test covered it.
Leave the body alone and report the assigned id on a new
CamelPulsarProducerMessageId header, named
after the producer headers this component already has. Documented in the
4.23 upgrade guide, since a
route reading the id from the body has to read the header instead.
Co-Authored-By: Claude Opus 5 <[email protected]>
Signed-off-by: Andrea Cosentino <[email protected]>
* CAMEL-24935: camel-pulsar - address review feedback
Rename the constant to PRODUCER_MESSAGE_ID. The other *_OUT constants in
PulsarMessageHeaders are values
the producer reads, while this one is a result it writes, so the suffix was
misleading. The header value
is unchanged; only the catalog constantName moves with it.
Note in the upgrade guide that the body assignment dates from the
asynchronous send added in 3.20
(CAMEL-16030).
Co-Authored-By: Claude Opus 5 <[email protected]>
Signed-off-by: Andrea Cosentino <[email protected]>
---------
Signed-off-by: Andrea Cosentino <[email protected]>
Co-authored-by: Claude Opus 5 <[email protected]>
---
.../apache/camel/catalog/components/pulsar.json | 3 +-
.../org/apache/camel/component/pulsar/pulsar.json | 3 +-
.../camel/component/pulsar/PulsarProducer.java | 2 +-
.../pulsar/utils/message/PulsarMessageHeaders.java | 3 +
.../pulsar/PulsarProducerMessageIdHeaderTest.java | 90 ++++++++++++++++++++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 12 +++
.../endpoint/dsl/PulsarEndpointBuilderFactory.java | 12 +++
7 files changed, 122 insertions(+), 3 deletions(-)
diff --git
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/components/pulsar.json
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/components/pulsar.json
index d2076b2f1f72..d01d9f55f0bf 100644
---
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/components/pulsar.json
+++
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/components/pulsar.json
@@ -91,7 +91,8 @@
"CamelPulsarProducerMessageEventTime": { "index": 12, "kind": "header",
"displayName": "", "group": "producer", "label": "producer", "required": false,
"javaType": "Long", "deprecated": false, "deprecationNote": "", "autowired":
false, "secret": false, "description": "The event time of the message
message.", "constantName":
"org.apache.camel.component.pulsar.utils.message.PulsarMessageHeaders#EVENT_TIME_OUT"
},
"CamelPulsarProducerMessageDeliverAt": { "index": 13, "kind": "header",
"displayName": "", "group": "producer", "label": "producer", "required": false,
"javaType": "Long", "deprecated": false, "deprecationNote": "", "autowired":
false, "secret": false, "description": "Deliver the message only at or after
the specified absolute timestamp. The timestamp is milliseconds and based on
UTC (eg: System.currentTimeMillis) Note: messages are only delivered with delay
when a consumer is consum [...]
"CamelPulsarRedeliveryCount": { "index": 14, "kind": "header",
"displayName": "", "group": "consumer", "label": "consumer", "required": false,
"javaType": "int", "deprecated": false, "deprecationNote": "", "autowired":
false, "secret": false, "description": "The message redelivery count,
redelivery count maintain in pulsar broker.", "constantName":
"org.apache.camel.component.pulsar.utils.message.PulsarMessageHeaders#PULSAR_REDELIVERY_COUNT"
},
- "CamelPulsarProducerMessageDeliverAfter": { "index": 15, "kind": "header",
"displayName": "", "group": "producer", "label": "producer", "required": false,
"javaType": "Long", "deprecated": false, "deprecationNote": "", "autowired":
false, "secret": false, "description": "Deliver the message after a given
delayed time (millis).", "constantName":
"org.apache.camel.component.pulsar.utils.message.PulsarMessageHeaders#DELIVER_AFTER"
}
+ "CamelPulsarProducerMessageDeliverAfter": { "index": 15, "kind": "header",
"displayName": "", "group": "producer", "label": "producer", "required": false,
"javaType": "Long", "deprecated": false, "deprecationNote": "", "autowired":
false, "secret": false, "description": "Deliver the message after a given
delayed time (millis).", "constantName":
"org.apache.camel.component.pulsar.utils.message.PulsarMessageHeaders#DELIVER_AFTER"
},
+ "CamelPulsarProducerMessageId": { "index": 16, "kind": "header",
"displayName": "", "group": "producer", "label": "producer", "required": false,
"javaType": "org.apache.pulsar.client.api.MessageId", "deprecated": false,
"deprecationNote": "", "autowired": false, "secret": false, "description": "The
message id the broker assigned to the published message.", "constantName":
"org.apache.camel.component.pulsar.utils.message.PulsarMessageHeaders#PRODUCER_MESSAGE_ID"
}
},
"properties": {
"persistence": { "index": 0, "kind": "path", "displayName": "Persistence",
"group": "common", "label": "", "required": true, "type": "enum", "javaType":
"java.lang.String", "enum": [ "persistent", "non-persistent" ], "deprecated":
false, "deprecationNote": "", "autowired": false, "secret": false,
"description": "Whether the topic is persistent or non-persistent" },
diff --git
a/components/camel-pulsar/src/generated/resources/META-INF/org/apache/camel/component/pulsar/pulsar.json
b/components/camel-pulsar/src/generated/resources/META-INF/org/apache/camel/component/pulsar/pulsar.json
index d2076b2f1f72..d01d9f55f0bf 100644
---
a/components/camel-pulsar/src/generated/resources/META-INF/org/apache/camel/component/pulsar/pulsar.json
+++
b/components/camel-pulsar/src/generated/resources/META-INF/org/apache/camel/component/pulsar/pulsar.json
@@ -91,7 +91,8 @@
"CamelPulsarProducerMessageEventTime": { "index": 12, "kind": "header",
"displayName": "", "group": "producer", "label": "producer", "required": false,
"javaType": "Long", "deprecated": false, "deprecationNote": "", "autowired":
false, "secret": false, "description": "The event time of the message
message.", "constantName":
"org.apache.camel.component.pulsar.utils.message.PulsarMessageHeaders#EVENT_TIME_OUT"
},
"CamelPulsarProducerMessageDeliverAt": { "index": 13, "kind": "header",
"displayName": "", "group": "producer", "label": "producer", "required": false,
"javaType": "Long", "deprecated": false, "deprecationNote": "", "autowired":
false, "secret": false, "description": "Deliver the message only at or after
the specified absolute timestamp. The timestamp is milliseconds and based on
UTC (eg: System.currentTimeMillis) Note: messages are only delivered with delay
when a consumer is consum [...]
"CamelPulsarRedeliveryCount": { "index": 14, "kind": "header",
"displayName": "", "group": "consumer", "label": "consumer", "required": false,
"javaType": "int", "deprecated": false, "deprecationNote": "", "autowired":
false, "secret": false, "description": "The message redelivery count,
redelivery count maintain in pulsar broker.", "constantName":
"org.apache.camel.component.pulsar.utils.message.PulsarMessageHeaders#PULSAR_REDELIVERY_COUNT"
},
- "CamelPulsarProducerMessageDeliverAfter": { "index": 15, "kind": "header",
"displayName": "", "group": "producer", "label": "producer", "required": false,
"javaType": "Long", "deprecated": false, "deprecationNote": "", "autowired":
false, "secret": false, "description": "Deliver the message after a given
delayed time (millis).", "constantName":
"org.apache.camel.component.pulsar.utils.message.PulsarMessageHeaders#DELIVER_AFTER"
}
+ "CamelPulsarProducerMessageDeliverAfter": { "index": 15, "kind": "header",
"displayName": "", "group": "producer", "label": "producer", "required": false,
"javaType": "Long", "deprecated": false, "deprecationNote": "", "autowired":
false, "secret": false, "description": "Deliver the message after a given
delayed time (millis).", "constantName":
"org.apache.camel.component.pulsar.utils.message.PulsarMessageHeaders#DELIVER_AFTER"
},
+ "CamelPulsarProducerMessageId": { "index": 16, "kind": "header",
"displayName": "", "group": "producer", "label": "producer", "required": false,
"javaType": "org.apache.pulsar.client.api.MessageId", "deprecated": false,
"deprecationNote": "", "autowired": false, "secret": false, "description": "The
message id the broker assigned to the published message.", "constantName":
"org.apache.camel.component.pulsar.utils.message.PulsarMessageHeaders#PRODUCER_MESSAGE_ID"
}
},
"properties": {
"persistence": { "index": 0, "kind": "path", "displayName": "Persistence",
"group": "common", "label": "", "required": true, "type": "enum", "javaType":
"java.lang.String", "enum": [ "persistent", "non-persistent" ], "deprecated":
false, "deprecationNote": "", "autowired": false, "secret": false,
"description": "Whether the topic is persistent or non-persistent" },
diff --git
a/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/PulsarProducer.java
b/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/PulsarProducer.java
index b83505d3b096..88840855aaf6 100644
---
a/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/PulsarProducer.java
+++
b/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/PulsarProducer.java
@@ -82,7 +82,7 @@ public class PulsarProducer extends DefaultAsyncProducer {
}
messageBuilder.sendAsync()
- .thenAccept(r -> exchange.getIn().setBody(r))
+ .thenAccept(r ->
exchange.getIn().setHeader(PulsarMessageHeaders.PRODUCER_MESSAGE_ID, r))
.whenComplete(
(r, e) -> {
try {
diff --git
a/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageHeaders.java
b/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageHeaders.java
index 9b5dbe298585..246b7dae519d 100644
---
a/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageHeaders.java
+++
b/components/camel-pulsar/src/main/java/org/apache/camel/component/pulsar/utils/message/PulsarMessageHeaders.java
@@ -60,4 +60,7 @@ public interface PulsarMessageHeaders {
@Metadata(label = "producer", description = "Deliver the message after a
given delayed time (millis).",
javaType = "Long")
String DELIVER_AFTER = "CamelPulsarProducerMessageDeliverAfter";
+ @Metadata(label = "producer", description = "The message id the broker
assigned to the published message.",
+ javaType = "org.apache.pulsar.client.api.MessageId")
+ String PRODUCER_MESSAGE_ID = "CamelPulsarProducerMessageId";
}
diff --git
a/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/PulsarProducerMessageIdHeaderTest.java
b/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/PulsarProducerMessageIdHeaderTest.java
new file mode 100644
index 000000000000..7f877c49894b
--- /dev/null
+++
b/components/camel-pulsar/src/test/java/org/apache/camel/component/pulsar/PulsarProducerMessageIdHeaderTest.java
@@ -0,0 +1,90 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.component.pulsar;
+
+import java.util.concurrent.CompletableFuture;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.Exchange;
+import org.apache.camel.RoutesBuilder;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.pulsar.utils.message.PulsarMessageHeaders;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.apache.pulsar.client.api.MessageId;
+import org.apache.pulsar.client.api.Producer;
+import org.apache.pulsar.client.api.ProducerBuilder;
+import org.apache.pulsar.client.api.PulsarClient;
+import org.apache.pulsar.client.api.TypedMessageBuilder;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.mockito.Answers.RETURNS_SELF;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+/**
+ * The send result used to replace the message body, which left the rest of
the route without the payload it had sent.
+ */
+public class PulsarProducerMessageIdHeaderTest extends CamelTestSupport {
+
+ private final MessageId messageId = mock(MessageId.class);
+
+ @Test
+ public void testTheBodySurvivesAndTheMessageIdIsAHeader() {
+ final Exchange out = template.request("direct:start", exchange ->
exchange.getIn().setBody("Hello World!"));
+
+ assertEquals("Hello World!", out.getMessage().getBody(String.class),
+ "the producer should leave the body alone");
+ assertSame(messageId,
out.getMessage().getHeader(PulsarMessageHeaders.PRODUCER_MESSAGE_ID),
+ "the send result should be reported as a header");
+ }
+
+ @Override
+ protected CamelContext createCamelContext() throws Exception {
+ final CamelContext context = super.createCamelContext();
+
+ final TypedMessageBuilder<byte[]> messageBuilder =
mock(TypedMessageBuilder.class, RETURNS_SELF);
+
when(messageBuilder.sendAsync()).thenReturn(CompletableFuture.completedFuture(messageId));
+
+ final Producer<byte[]> pulsarProducer = mock(Producer.class);
+ when(pulsarProducer.newMessage()).thenReturn(messageBuilder);
+
+ final ProducerBuilder<byte[]> producerBuilder =
mock(ProducerBuilder.class, RETURNS_SELF);
+ when(producerBuilder.create()).thenReturn(pulsarProducer);
+
+ final PulsarClient pulsarClient = mock(PulsarClient.class);
+ when(pulsarClient.newProducer()).thenReturn(producerBuilder);
+
+ final PulsarComponent component = new PulsarComponent(context);
+ component.setPulsarClient(pulsarClient);
+ context.addComponent("pulsar", component);
+
+ return context;
+ }
+
+ @Override
+ protected RoutesBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start")
+
.to("pulsar:persistent://public/default/camel-producer-test");
+ }
+ };
+ }
+}
diff --git
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
index bfdb4e9a5260..3bcd2615c9e2 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
@@ -2746,6 +2746,18 @@ was published as a second, undocumented spelling of the
same setting. Routes tha
`databaseOutServerName` should use `xstreamOutServerName` instead; it
configures the same XStream outbound
server and is unchanged.
+=== camel-pulsar - the producer no longer replaces the body with the message id
+
+After sending, the producer used to overwrite the message body with the
`MessageId` returned by the
+broker, so every step after the `to("pulsar:...")` saw a `MessageId` instead
of the payload, and a route
+such as `.to("pulsar:a").to("pulsar:b")` published a Java-serialized
`MessageId` to the second topic.
+That behaviour dates from the asynchronous send added in 3.20 (CAMEL-16030).
+
+The body is now left untouched, and the assigned id is reported on the new
+`CamelPulsarProducerMessageId` header instead.
+
+A route that read the `MessageId` from the body must read that header instead.
+
=== camel-jolt - switch from unmaintained bazaarvoice jolt to the
jolt-community fork
The JOLT library dependency has been migrated from
`com.bazaarvoice.jolt:jolt-core` to
diff --git
a/dsl/camel-endpointdsl/src/generated/java/org/apache/camel/builder/endpoint/dsl/PulsarEndpointBuilderFactory.java
b/dsl/camel-endpointdsl/src/generated/java/org/apache/camel/builder/endpoint/dsl/PulsarEndpointBuilderFactory.java
index 5402c8322e8c..3a73d95d10c0 100644
---
a/dsl/camel-endpointdsl/src/generated/java/org/apache/camel/builder/endpoint/dsl/PulsarEndpointBuilderFactory.java
+++
b/dsl/camel-endpointdsl/src/generated/java/org/apache/camel/builder/endpoint/dsl/PulsarEndpointBuilderFactory.java
@@ -1849,6 +1849,18 @@ public interface PulsarEndpointBuilderFactory {
public String pulsarProducerMessageDeliverAfter() {
return "CamelPulsarProducerMessageDeliverAfter";
}
+ /**
+ * The message id the broker assigned to the published message.
+ *
+ * The option is a: {@code org.apache.pulsar.client.api.MessageId}
type.
+ *
+ * Group: producer
+ *
+ * @return the name of the header {@code PulsarProducerMessageId}.
+ */
+ public String pulsarProducerMessageId() {
+ return "CamelPulsarProducerMessageId";
+ }
}
static PulsarEndpointBuilder endpointBuilder(String componentName, String
path) {
class PulsarEndpointBuilderImpl extends AbstractEndpointBuilder
implements PulsarEndpointBuilder, AdvancedPulsarEndpointBuilder {