From 650d927feb630e03f37352ba181c6700859185d5 Mon Sep 17 00:00:00 2001 From: Andrey Yegorov Date: Wed, 28 Aug 2024 12:29:29 -0700 Subject: [PATCH] trying to repro issue with null in udts --- .../sink/simulacron/AvroLogicalTypesTest.java | 2 +- .../oss/pulsar/sink/simulacron/AvroTest.java | 51 +++++++++++++++---- .../sink/simulacron/GrabNarFilePathTest.java | 2 +- .../sink/simulacron/PulsarCCMTestBase.java | 2 +- 4 files changed, 43 insertions(+), 14 deletions(-) diff --git a/tests/src/test/java/com/datastax/oss/pulsar/sink/simulacron/AvroLogicalTypesTest.java b/tests/src/test/java/com/datastax/oss/pulsar/sink/simulacron/AvroLogicalTypesTest.java index 455f5ae..4b79dd2 100644 --- a/tests/src/test/java/com/datastax/oss/pulsar/sink/simulacron/AvroLogicalTypesTest.java +++ b/tests/src/test/java/com/datastax/oss/pulsar/sink/simulacron/AvroLogicalTypesTest.java @@ -512,4 +512,4 @@ public static byte[] serializeAvroGenericRecord( throw new RuntimeException(e); } } -} \ No newline at end of file +} diff --git a/tests/src/test/java/com/datastax/oss/pulsar/sink/simulacron/AvroTest.java b/tests/src/test/java/com/datastax/oss/pulsar/sink/simulacron/AvroTest.java index 9038303..e86273f 100644 --- a/tests/src/test/java/com/datastax/oss/pulsar/sink/simulacron/AvroTest.java +++ b/tests/src/test/java/com/datastax/oss/pulsar/sink/simulacron/AvroTest.java @@ -59,6 +59,15 @@ public class AvroTest extends PulsarCCMTestBase { ImmutableMap.of("k1", 7.0D, "k2", 9.0D)); private final List listOfUdt = ImmutableList.of(pojoUdt, pojoUdt); + private final MyUdt pojoUdtWithNull = + new MyUdt( + 99, + null, + ImmutableList.of("l1", "l2"), + ImmutableSet.of(3, 4), + ImmutableMap.of("k1", 7.0D, "k2", 9.0D)); + private final List listOfUdtWithNull = ImmutableList.of(pojoUdtWithNull, pojoUdtWithNull); + /** * AVRO schema with Pulsar doesn't work well with mixed value types on the map - the values will * be of "org.apache.avro.generic.GenericData$Record" with the following limitations: 1. Using @@ -110,6 +119,26 @@ protected void performTest(final PulsarSinkTester pulsarSink) throws PulsarClien mapUdt, listOfUdt)) .send(); + producer + .newMessage() + .key("838") + .value( + new MyBean( + "value1", + map, + list, + set, + listOfMaps, + setOfMaps, + mapOfLists, + mapOfSets, + listOfSets, + setOfLists, + pojoUdtWithNull, + null, + mapUdt, + listOfUdtWithNull)) + .send(); } try { Awaitility.waitAtMost(30, TimeUnit.SECONDS) @@ -117,7 +146,7 @@ protected void performTest(final PulsarSinkTester pulsarSink) throws PulsarClien .until( () -> { List results = session.execute("SELECT * FROM table1").all(); - return results.size() > 0; + return results.size() > 0 && null == results.get(0).getUdtValue("f").getString("stringf"); }); List results = session.execute("SELECT * FROM table1").all(); @@ -156,14 +185,14 @@ protected void performTest(final PulsarSinkTester pulsarSink) throws PulsarClien DefaultUdtValue value = (DefaultUdtValue) row.getUdtValue("f"); assertEquals(value.size(), 5); - assertEquals(pojoUdt.getIntf(), value.getInt("intf")); - assertEquals(pojoUdt.getStringf(), value.getString("stringf")); + assertEquals(pojoUdtWithNull.getIntf(), value.getInt("intf")); + assertEquals(pojoUdtWithNull.getStringf(), value.getString("stringf")); GenericType> udtListType = new GenericType>() {}; - assertEquals(pojoUdt.getListf(), value.get("listf", udtListType)); + assertEquals(pojoUdtWithNull.getListf(), value.get("listf", udtListType)); GenericType> udtSetType = new GenericType>() {}; - assertEquals(pojoUdt.getSetf(), value.get("setf", udtSetType)); + assertEquals(pojoUdtWithNull.getSetf(), value.get("setf", udtSetType)); GenericType> udtMapType = new GenericType>() {}; - assertEquals(pojoUdt.getMapf(), value.get("mapf", udtMapType)); + assertEquals(pojoUdtWithNull.getMapf(), value.get("mapf", udtMapType)); value = (DefaultUdtValue) row.getUdtValue("g"); assertEquals(Integer.valueOf(mapUdt.get("intf").toString()), value.getInt("intf")); @@ -174,11 +203,11 @@ protected void performTest(final PulsarSinkTester pulsarSink) throws PulsarClien List listOfUdt = row.get("s", listOfUdtType); assertEquals(listOfUdt.size(), 2); for (UdtValue udt : listOfUdt) { - assertEquals(pojoUdt.getIntf(), udt.getInt("intf")); - assertEquals(pojoUdt.getStringf(), udt.getString("stringf")); - assertEquals(pojoUdt.getListf(), udt.get("listf", udtListType)); - assertEquals(pojoUdt.getSetf(), udt.get("setf", udtSetType)); - assertEquals(pojoUdt.getMapf(), udt.get("mapf", udtMapType)); + assertEquals(pojoUdtWithNull.getIntf(), udt.getInt("intf")); + assertEquals(pojoUdtWithNull.getStringf(), udt.getString("stringf")); + assertEquals(pojoUdtWithNull.getListf(), udt.get("listf", udtListType)); + assertEquals(pojoUdtWithNull.getSetf(), udt.get("setf", udtSetType)); + assertEquals(pojoUdtWithNull.getMapf(), udt.get("mapf", udtMapType)); } } assertEquals(1, results.size()); diff --git a/tests/src/test/java/com/datastax/oss/pulsar/sink/simulacron/GrabNarFilePathTest.java b/tests/src/test/java/com/datastax/oss/pulsar/sink/simulacron/GrabNarFilePathTest.java index 9b98c46..7b8930a 100644 --- a/tests/src/test/java/com/datastax/oss/pulsar/sink/simulacron/GrabNarFilePathTest.java +++ b/tests/src/test/java/com/datastax/oss/pulsar/sink/simulacron/GrabNarFilePathTest.java @@ -48,4 +48,4 @@ public void test() { assertThat(narPath, containsString(".nar")); assertTrue(new File(narPath).isFile()); } -} \ No newline at end of file +} diff --git a/tests/src/test/java/com/datastax/oss/pulsar/sink/simulacron/PulsarCCMTestBase.java b/tests/src/test/java/com/datastax/oss/pulsar/sink/simulacron/PulsarCCMTestBase.java index 73c24b2..3eacd31 100644 --- a/tests/src/test/java/com/datastax/oss/pulsar/sink/simulacron/PulsarCCMTestBase.java +++ b/tests/src/test/java/com/datastax/oss/pulsar/sink/simulacron/PulsarCCMTestBase.java @@ -395,4 +395,4 @@ public void setListOfUdt(List listOfUdt) { this.listOfUdt = listOfUdt; } } -} \ No newline at end of file +}