From c75a2832672665c9986cab74e498db2c9cf56799 Mon Sep 17 00:00:00 2001 From: Nicolas Gibanel Date: Sun, 13 Sep 2026 13:04:07 +0200 Subject: [PATCH] SolaceIO: map message user properties --- CHANGES.md | 1 + .../beam/sdk/io/solace/data/Solace.java | 68 ++++++++++++++++- .../solace/data/SolaceRecordMapperTest.java | 75 +++++++++++++++++++ .../sdk/io/solace/data/SolaceRecordTest.java | 21 ++++++ 4 files changed, 164 insertions(+), 1 deletion(-) 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();