diff --git a/CHANGES.md b/CHANGES.md
index f0d5d06b9d9f..4f2edcf270f5 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -69,6 +69,7 @@
* BigQueryIO now supports reading BigQuery Lakehouse runtime catalog (BigLake metastore) Iceberg tables with the Storage Read API, using 4-part `project.catalog.namespace.table` identifiers (or a `TableReference` with a composite `catalog.namespace` dataset id). Previously such references were silently mis-parsed (Java) ([#39597](https://github.com/apache/beam/issues/39597)) .
* SolaceIO now supports reading and writing binary and text content data payload (Java) ([#39875](https://github.com/apache/beam/issues/39875)).
* ClickHouseIO: support writing `Decimal(P, S)` / `Decimal32/64/128/256` columns (Java) ([#39840](https://github.com/apache/beam/issues/39840)).
+* SolaceIO now supports reading and writing user properties (message metadata) (Java) ([#40099](https://github.com/apache/beam/issues/40099)).
## New Features / Improvements
diff --git a/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/data/Solace.java b/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/data/Solace.java
index 15fe06103fbd..b9997fec4229 100644
--- a/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/data/Solace.java
+++ b/sdks/java/io/solace/src/main/java/org/apache/beam/sdk/io/solace/data/Solace.java
@@ -21,12 +21,17 @@
import com.solacesystems.jcsmp.BytesMessage;
import com.solacesystems.jcsmp.BytesXMLMessage;
import com.solacesystems.jcsmp.JCSMPFactory;
+import com.solacesystems.jcsmp.SDTException;
+import com.solacesystems.jcsmp.SDTMap;
import com.solacesystems.jcsmp.TextMessage;
import java.nio.ByteBuffer;
import java.nio.charset.CharacterCodingException;
import java.nio.charset.CodingErrorAction;
import java.nio.charset.StandardCharsets;
import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
import org.apache.beam.sdk.schemas.AutoValueSchema;
import org.apache.beam.sdk.schemas.annotations.DefaultSchema;
import org.apache.beam.sdk.schemas.annotations.SchemaFieldNumber;
@@ -276,6 +281,16 @@ public enum PayloadType {
@SchemaFieldNumber("13")
public abstract PayloadType getPayloadType();
+ /**
+ * Gets the user properties of the message as a string map.
+ *
+ *
Mapped from {@link BytesXMLMessage#getProperties()}. Values are stringified.
+ *
+ * @return The user properties, or an empty map if the message carries none.
+ */
+ @SchemaFieldNumber("14")
+ public abstract Map getUserProperties();
+
/** Gets the payload decoded as UTF-8 when this record has type {@link PayloadType#TEXT}. */
public final String getText() {
if (getPayloadType() != PayloadType.TEXT) {
@@ -292,7 +307,8 @@ public static Builder builder() {
.setRedelivered(false)
.setTimeToLive(0)
.setAttachmentBytes(new byte[0])
- .setPayloadType(PayloadType.BYTES_XML);
+ .setPayloadType(PayloadType.BYTES_XML)
+ .setUserProperties(Collections.emptyMap());
}
@AutoValue.Builder
@@ -332,6 +348,8 @@ public abstract Builder setReplicationGroupMessageId(
public abstract Builder setAttachmentBytes(byte[] attachmentBytes);
+ public abstract Builder setUserProperties(Map userProperties);
+
public abstract Record build();
}
@@ -456,6 +474,7 @@ public static class SolaceRecordMapper {
Destination replyTo = getDestination(msg.getCorrelationId(), msg.getReplyTo());
Destination destination = getDestination(msg.getCorrelationId(), msg.getDestination());
+ Map userProperties = getUserProperties(msg.getProperties());
Record.Builder recordBuilder = decodePayload(msg);
return recordBuilder
@@ -473,6 +492,7 @@ public static class SolaceRecordMapper {
msg.getReplicationGroupMessageId() != null
? msg.getReplicationGroupMessageId().toString()
: null)
+ .setUserProperties(userProperties)
.build();
}
@@ -519,6 +539,10 @@ public static BytesXMLMessage toMessage(Record record) {
msg.setSenderTimestamp(senderTimestamp);
msg.setApplicationMessageId(record.getMessageId());
+ if (!record.getUserProperties().isEmpty()) {
+ msg.setProperties(createUserProperties(record.getUserProperties()));
+ }
+
return msg;
}
@@ -599,5 +623,47 @@ private static byte[] readAttachment(BytesXMLMessage msg) {
buffer.get(attachment);
return attachment;
}
+
+ private static Map getUserProperties(@Nullable SDTMap properties) {
+ if (properties == null || properties.isEmpty()) {
+ return Collections.emptyMap();
+ }
+
+ Map userProperties = new HashMap<>();
+ for (String key : properties.keySet()) {
+ String value = stringifyUserProperty(properties, key);
+ if (value == null) {
+ LOG.warn("User property '{}' has a null value, skipping.", key);
+ continue;
+ }
+ userProperties.put(key, value);
+ }
+ return Collections.unmodifiableMap(userProperties);
+ }
+
+ private static @Nullable String stringifyUserProperty(SDTMap properties, String key) {
+ try {
+ Object value = properties.get(key);
+ if (value == null) {
+ return null;
+ }
+ return String.valueOf(value);
+ } catch (SDTException e) {
+ LOG.error("Could not read user property '{}'.", key, e);
+ return null;
+ }
+ }
+
+ private static SDTMap createUserProperties(Map userProperties) {
+ SDTMap properties = JCSMPFactory.onlyInstance().createMap();
+ for (Map.Entry entry : userProperties.entrySet()) {
+ try {
+ properties.putString(entry.getKey(), entry.getValue());
+ } catch (SDTException e) {
+ LOG.error("Could not write user property '{}'.", entry.getKey(), e);
+ }
+ }
+ return properties;
+ }
}
}
diff --git a/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/data/SolaceRecordMapperTest.java b/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/data/SolaceRecordMapperTest.java
index cc6567b4c88c..6bf3c7248195 100644
--- a/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/data/SolaceRecordMapperTest.java
+++ b/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/data/SolaceRecordMapperTest.java
@@ -26,9 +26,13 @@
import com.solacesystems.jcsmp.BytesXMLMessage;
import com.solacesystems.jcsmp.DeliveryMode;
import com.solacesystems.jcsmp.JCSMPFactory;
+import com.solacesystems.jcsmp.SDTMap;
import com.solacesystems.jcsmp.TextMessage;
import java.nio.charset.StandardCharsets;
import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
import org.apache.beam.sdk.io.solace.broker.MessageProducerUtils;
import org.apache.beam.sdk.io.solace.data.Solace.Record;
import org.apache.beam.sdk.io.solace.data.Solace.Record.PayloadType;
@@ -144,6 +148,34 @@ public void testMapMessageMetadata() {
assertEquals(789L, record.getTimeToLive());
}
+ @Test
+ public void testMapMessageUserProperties() throws Exception {
+ BytesXMLMessage message = JCSMPFactory.onlyInstance().createBytesXMLMessage();
+ message.setApplicationMessageId("id");
+ SDTMap properties = JCSMPFactory.onlyInstance().createMap();
+ properties.putString("contentType", "application/json");
+ properties.putInteger("attempt", 3);
+ properties.putString("null", null);
+ message.setProperties(properties);
+
+ Record record = Solace.SolaceRecordMapper.toRecord(message);
+
+ Map expected = new HashMap<>();
+ expected.put("contentType", "application/json");
+ expected.put("attempt", "3");
+ assertEquals(expected, record.getUserProperties());
+ }
+
+ @Test
+ public void testMapWithEmptyMessageUserProperties() {
+ BytesXMLMessage message = JCSMPFactory.onlyInstance().createBytesXMLMessage();
+ message.setApplicationMessageId("id");
+
+ Record record = Solace.SolaceRecordMapper.toRecord(message);
+
+ assertTrue(record.getUserProperties().isEmpty());
+ }
+
@Test
public void testMapTextRecord() {
Record record =
@@ -257,9 +289,52 @@ public void testToMessageDoesNotSetPublishingFields() {
assertNull(msg.getCorrelationKey());
}
+ @Test
+ public void testMapRecordUserProperties() throws Exception {
+ Record record =
+ Record.builder()
+ .setMessageId("id")
+ .setText("hello")
+ .setSenderTimestamp(1L)
+ .setUserProperties(Collections.singletonMap("contentType", "application/json"))
+ .build();
+
+ BytesXMLMessage msg = Solace.SolaceRecordMapper.toMessage(record);
+
+ assertEquals("application/json", msg.getProperties().getString("contentType"));
+ }
+
+ @Test
+ public void testMapWithEmptyRecordUserProperties() {
+ Record record =
+ Record.builder().setMessageId("id").setText("hello").setSenderTimestamp(1L).build();
+
+ BytesXMLMessage msg = Solace.SolaceRecordMapper.toMessage(record);
+
+ assertNull(msg.getProperties());
+ }
+
// ---------------------------------------------------------------------------
// round-trip
// ---------------------------------------------------------------------------
+ @Test
+ public void testRoundTripUserProperties() {
+ Record original =
+ Record.builder()
+ .setMessageId("id")
+ .setText("hello")
+ .setSenderTimestamp(1L)
+ .setUserProperties(Collections.singletonMap("contentType", "application/json"))
+ .build();
+
+ BytesXMLMessage msg = Solace.SolaceRecordMapper.toMessage(original);
+ msg.setApplicationMessageId("id");
+ Record decoded = Solace.SolaceRecordMapper.toRecord(msg);
+
+ assertEquals(
+ Collections.singletonMap("contentType", "application/json"), decoded.getUserProperties());
+ }
+
@Test
public void testRoundTripTextPayload() {
Record original =
diff --git a/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/data/SolaceRecordTest.java b/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/data/SolaceRecordTest.java
index 5b521f7d4e68..0e8946bde472 100644
--- a/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/data/SolaceRecordTest.java
+++ b/sdks/java/io/solace/src/test/java/org/apache/beam/sdk/io/solace/data/SolaceRecordTest.java
@@ -19,8 +19,10 @@
import static org.junit.Assert.assertArrayEquals;
import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertTrue;
import java.nio.charset.StandardCharsets;
+import java.util.Collections;
import org.apache.beam.sdk.io.solace.data.Solace.Record;
import org.junit.Test;
@@ -33,6 +35,25 @@ public void testDefaultPayloadType() {
assertEquals(Record.PayloadType.BYTES_XML, record.getPayloadType());
}
+ @Test
+ public void testDefaultUserPropertiesIsEmpty() {
+ Record record = Record.builder().setMessageId("id").setPayload(new byte[0]).build();
+
+ assertTrue(record.getUserProperties().isEmpty());
+ }
+
+ @Test
+ public void testSetUserProperties() {
+ Record record =
+ Record.builder()
+ .setMessageId("id")
+ .setPayload(new byte[0])
+ .setUserProperties(Collections.singletonMap("key", "value"))
+ .build();
+
+ assertEquals(Collections.singletonMap("key", "value"), record.getUserProperties());
+ }
+
@Test
public void testSetTextPayload() {
Record record = Record.builder().setMessageId("id").setText("héllo").build();