From 609a1f90456115932a2f873698317746f3be3c50 Mon Sep 17 00:00:00 2001 From: david-streamlio Date: Sat, 18 Jun 2022 19:09:47 -0700 Subject: [PATCH 01/14] Added GenericRecord support --- pom.xml | 6 +- pulsar-io/cassandra/pom.xml | 21 +++ .../io/cassandra/CassandraAbstractSink.java | 54 +++----- .../io/cassandra/CassandraConnector.java | 130 ++++++++++++++++++ .../pulsar/io/cassandra/CassandraSink.java | 104 ++++++++++++++ .../io/cassandra/CassandraSinkConfig.java | 20 +++ .../io/cassandra/CassandraStringSink.java | 2 +- .../META-INF/services/pulsar-io.yaml | 2 +- .../io/cassandra/CassandraConnectorTest.java | 55 ++++++++ .../io/cassandra/CassandraSinkExec.java | 98 +++++++++++++ .../AbstractGenericRecordProducer.java | 43 ++++++ .../producers/InputTopicProducerThread.java | 83 +++++++++++ .../producers/InputTopicStringProducer.java | 43 ++++++ .../ObservationSchemaRecordProducer.java | 60 ++++++++ .../ReadingSchemaRecordProducer.java | 93 +++++++++++++ .../test/resources/cassandra-sink-config.yaml | 23 ++++ tests/integration/pom.xml | 3 +- 17 files changed, 796 insertions(+), 44 deletions(-) create mode 100644 pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraConnector.java create mode 100644 pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSink.java create mode 100644 pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraConnectorTest.java create mode 100644 pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java create mode 100644 pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/AbstractGenericRecordProducer.java create mode 100644 pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicProducerThread.java create mode 100644 pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java create mode 100644 pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ObservationSchemaRecordProducer.java create mode 100644 pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ReadingSchemaRecordProducer.java create mode 100644 pulsar-io/cassandra/src/test/resources/cassandra-sink-config.yaml diff --git a/pom.xml b/pom.xml index ac5d81ea008ca..b57c3475e1e74 100644 --- a/pom.xml +++ b/pom.xml @@ -1078,11 +1078,7 @@ flexible messaging model and an intuitive client API. pom import - - com.datastax.cassandra - cassandra-driver-core - ${cassandra.version} - + org.assertj assertj-core diff --git a/pulsar-io/cassandra/pom.xml b/pulsar-io/cassandra/pom.xml index cbeeb1d272b24..abf18b9c638bc 100644 --- a/pulsar-io/cassandra/pom.xml +++ b/pulsar-io/cassandra/pom.xml @@ -51,6 +51,27 @@ com.datastax.cassandra cassandra-driver-core + ${cassandra.version} + + + io.dropwizard.metrics + metrics-core + + + + + + org.apache.pulsar + pulsar-functions-local-runner-original + 2.11.0-SNAPSHOT + test + + + + org.apache.pulsar + pulsar-io-common + 2.11.0-SNAPSHOT + compile diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraAbstractSink.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraAbstractSink.java index 7a872ff6504ec..1323379c66c9b 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraAbstractSink.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraAbstractSink.java @@ -19,17 +19,14 @@ package org.apache.pulsar.io.cassandra; -import com.datastax.driver.core.BoundStatement; -import com.datastax.driver.core.Cluster; -import com.datastax.driver.core.PreparedStatement; +import com.datastax.driver.core.*; import com.datastax.driver.core.ResultSet; -import com.datastax.driver.core.ResultSetFuture; -import com.datastax.driver.core.Session; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.MoreExecutors; import java.util.Map; import org.apache.pulsar.functions.api.Record; +import org.apache.pulsar.io.common.IOConfigUtils; import org.apache.pulsar.io.core.KeyValue; import org.apache.pulsar.io.core.Sink; import org.apache.pulsar.io.core.SinkContext; @@ -40,15 +37,14 @@ */ public abstract class CassandraAbstractSink implements Sink { - // ----- Runtime fields - private Cluster cluster; - private Session session; - CassandraSinkConfig cassandraSinkConfig; - private PreparedStatement statement; + private CassandraConnector connector; + private CassandraSinkConfig cassandraSinkConfig; @Override public void open(Map config, SinkContext sinkContext) throws Exception { - cassandraSinkConfig = CassandraSinkConfig.load(config); + + cassandraSinkConfig = IOConfigUtils.loadWithSecrets(config, CassandraSinkConfig.class, sinkContext); + if (cassandraSinkConfig.getRoots() == null || cassandraSinkConfig.getKeyspace() == null || cassandraSinkConfig.getKeyname() == null @@ -56,22 +52,26 @@ public void open(Map config, SinkContext sinkContext) throws Exc || cassandraSinkConfig.getColumnName() == null) { throw new IllegalArgumentException("Required property not set."); } - createClient(cassandraSinkConfig.getRoots()); - statement = session.prepare("INSERT INTO " + cassandraSinkConfig.getColumnFamily() + " (" - + cassandraSinkConfig.getKeyname() + ", " + cassandraSinkConfig.getColumnName() + ") VALUES (?, ?)"); + + connector = new CassandraConnector(cassandraSinkConfig); + connector.connect(); } @Override public void close() throws Exception { - session.close(); - cluster.close(); + connector.close(); } @Override public void write(Record record) { + KeyValue keyValue = extractKeyValue(record); - BoundStatement bound = statement.bind(keyValue.getKey(), keyValue.getValue()); - ResultSetFuture future = session.executeAsync(bound); + + BoundStatement bound = connector.getPreparedStatement() + .bind(keyValue.getKey(), keyValue.getValue()); + + ResultSetFuture future = connector.getSession().executeAsync(bound); + Futures.addCallback(future, new FutureCallback() { @Override @@ -86,23 +86,5 @@ public void onFailure(Throwable t) { }, MoreExecutors.directExecutor()); } - private void createClient(String roots) { - String[] hosts = roots.split(","); - if (hosts.length <= 0) { - throw new RuntimeException("Invalid cassandra roots"); - } - Cluster.Builder b = Cluster.builder(); - for (int i = 0; i < hosts.length; ++i) { - String[] hostPort = hosts[i].split(":"); - b.addContactPoint(hostPort[0]); - if (hostPort.length > 1) { - b.withPort(Integer.parseInt(hostPort[1])); - } - } - cluster = b.build(); - session = cluster.connect(); - session.execute("USE " + cassandraSinkConfig.getKeyspace()); - } - public abstract KeyValue extractKeyValue(Record record); } \ No newline at end of file diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraConnector.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraConnector.java new file mode 100644 index 0000000000000..b4dede9eaf053 --- /dev/null +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraConnector.java @@ -0,0 +1,130 @@ +/** + * 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.pulsar.io.cassandra; + +import com.datastax.driver.core.*; +import com.datastax.driver.core.Session; + +import java.util.ArrayList; +import java.util.List; + +public class CassandraConnector { + + private Cluster cluster; + private Session session; + private PreparedStatement statement; + private List tableFields; + private final CassandraSinkConfig config; + + public CassandraConnector(CassandraSinkConfig config) { + this.config = config; + } + + public void connect() { + session = getCluster().connect(config.getKeyspace()); + } + + public Session getSession() { + if (session == null) { + this.connect(); + } + return session; + } + + public PreparedStatement getPreparedStatement() { + if (statement == null) { + List fields = getTableFields(); + + StringBuilder sb = new StringBuilder("INSERT INTO ") + .append(config.getKeyspace() + "." + config.getColumnFamily() + " ("); + + for (int idx = 0; idx < fields.size(); idx++) { + sb.append(fields.get(idx)); + if (idx < fields.size()-1) { + sb.append(", "); + } + } + + sb.append(") VALUES ("); + + for (int idx = 0; idx < fields.size(); idx++) { + sb.append("?"); + if (idx < fields.size()-1) { + sb.append(", "); + } + } + + sb.append(")"); + statement = getSession().prepare(sb.toString()); + } + + return statement; + } + + List getTableFields() { + + if (tableFields == null) { + + TableMetadata meta = getCluster().getMetadata() + .getKeyspace(config.getKeyspace()) + .getTable(config.getColumnFamily()); + + tableFields = new ArrayList(meta.getColumns().size()); + + for (ColumnMetadata col : meta.getColumns()) { + tableFields.add(col.getName()); + } + } + return tableFields; + } + + private Cluster getCluster() { + if (cluster == null) { + String[] hosts = config.getRoots().split(","); + + Cluster.Builder builder = Cluster.builder().withoutMetrics(); + + for (int i = 0; i < hosts.length; ++i) { + String[] hostPort = hosts[i].split(":"); + builder.addContactPoint(hostPort[0]); + if (hostPort.length > 1) { + builder.withPort(Integer.parseInt(hostPort[1])); + } + } + + // Authenticate if credentials have been provided + if (config.getUserName() != null + && config.getPassword() != null) { + builder.withCredentials( + config.getUserName(), + config.getPassword() + ); + } + cluster = builder.build(); + } + + return cluster; + } + + public void close() { + session.close(); + cluster.close(); + } +} diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSink.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSink.java new file mode 100644 index 0000000000000..470d5a0267715 --- /dev/null +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSink.java @@ -0,0 +1,104 @@ +/** + * 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.pulsar.io.cassandra; + +import com.datastax.driver.core.*; +import com.google.common.util.concurrent.FutureCallback; +import com.google.common.util.concurrent.Futures; +import com.google.common.util.concurrent.MoreExecutors; +import org.apache.pulsar.client.api.schema.GenericRecord; +import org.apache.pulsar.functions.api.Record; +import org.apache.pulsar.io.common.IOConfigUtils; +import org.apache.pulsar.io.core.Sink; +import org.apache.pulsar.io.core.SinkContext; + +import java.util.Map; + +public class CassandraSink implements Sink { + + private CassandraConnector connector; + private CassandraSinkConfig cassandraSinkConfig; + private PreparedStatement stmt; + private BoundStatement boundStatement; + + @Override + public void open(Map config, SinkContext ctx) throws Exception { + + cassandraSinkConfig = IOConfigUtils.loadWithSecrets(config, CassandraSinkConfig.class, ctx); + + if (cassandraSinkConfig.getRoots() == null + || cassandraSinkConfig.getKeyspace() == null + || cassandraSinkConfig.getColumnFamily() == null) { + throw new IllegalArgumentException("Required property not set."); + } + + connector = new CassandraConnector(cassandraSinkConfig); + connector.connect(); + } + + @Override + public void write(Record record) throws Exception { + + Object[] boundValues = new Object[connector.getTableFields().size()]; + GenericRecord generic = record.getValue(); + + for (int idx = 0; idx < connector.getTableFields().size(); idx++) { + String fieldName = connector.getTableFields().get(idx); + boundValues[idx] = generic.getField(fieldName); + } + + BoundStatement bs = getBoundStatement().bind(boundValues); + + ResultSetFuture future = connector.getSession().executeAsync(bs); + + Futures.addCallback(future, + new FutureCallback() { + @Override + public void onSuccess(ResultSet result) { + record.ack(); + } + + @Override + public void onFailure(Throwable t) { + record.fail(); + } + }, MoreExecutors.directExecutor()); + } + + @Override + public void close() throws Exception { + connector.close(); + } + + private PreparedStatement getStatement() { + if (stmt == null) { + stmt = connector.getPreparedStatement(); + } + return stmt; + } + + private BoundStatement getBoundStatement() { + if (boundStatement == null) { + boundStatement = getStatement().bind(); + boundStatement.setConsistencyLevel(ConsistencyLevel.ALL); + } + return boundStatement; + } +} diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSinkConfig.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSinkConfig.java index 821e973dec656..c9b52e225284a 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSinkConfig.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSinkConfig.java @@ -35,26 +35,46 @@ public class CassandraSinkConfig implements Serializable { private static final long serialVersionUID = 1L; + @FieldDoc( + required = false, + defaultValue = "", + sensitive = true, + help = "Username used to connect to the database specified by `jdbcUrl`" + ) + private String userName; + + @FieldDoc( + required = false, + defaultValue = "", + sensitive = true, + help = "Password used to connect to the database specified by `jdbcUrl`" + ) + private String password; + @FieldDoc( required = true, defaultValue = "", help = "A comma-separated list of cassandra hosts to connect to") private String roots; + @FieldDoc( required = true, defaultValue = "", help = "The key space used for writing pulsar messages to") private String keyspace; + @FieldDoc( required = true, defaultValue = "", help = "The key name of the cassandra column family") private String keyname; + @FieldDoc( required = true, defaultValue = "", help = "The cassandra column family name") private String columnFamily; + @FieldDoc( required = true, defaultValue = "", diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraStringSink.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraStringSink.java index 694b79b055be6..3ff3af732e8d1 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraStringSink.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraStringSink.java @@ -26,7 +26,7 @@ /** * Cassandra sink that treats incoming messages on the input topic as Strings - * and write identical key/value pairs. + * and writes key/value pairs. */ @Connector( name = "cassandra", diff --git a/pulsar-io/cassandra/src/main/resources/META-INF/services/pulsar-io.yaml b/pulsar-io/cassandra/src/main/resources/META-INF/services/pulsar-io.yaml index b4863210c55a0..8954d1621297a 100644 --- a/pulsar-io/cassandra/src/main/resources/META-INF/services/pulsar-io.yaml +++ b/pulsar-io/cassandra/src/main/resources/META-INF/services/pulsar-io.yaml @@ -19,5 +19,5 @@ name: cassandra description: Writes data into Cassandra -sinkClass: org.apache.pulsar.io.cassandra.CassandraStringSink +sinkClass: org.apache.pulsar.io.cassandra.CassandraSink sinkConfigClass: org.apache.pulsar.io.cassandra.CassandraSinkConfig diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraConnectorTest.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraConnectorTest.java new file mode 100644 index 0000000000000..91c063a9ee2e0 --- /dev/null +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraConnectorTest.java @@ -0,0 +1,55 @@ +package org.apache.pulsar.io.cassandra; + +import org.junit.Ignore; +import org.junit.Test; + +import static org.junit.Assert.assertNotNull; +import static org.testng.AssertJUnit.assertEquals; + +public class CassandraConnectorTest { + + private CassandraSinkConfig config; + + @Test + @Ignore + public final void securedTest() { + config = new CassandraSinkConfig(); + config.setRoots("localhost"); + config.setUserName("cassandra"); + config.setPassword("cassandra"); + + CassandraConnector connector = new CassandraConnector(config); + connector.connect(); + assertNotNull(connector.getSession()); + } + + @Test + public final void getObservationPreparedStatementTest() { + config = new CassandraSinkConfig(); + config.setRoots("localhost"); + config.setUserName("cassandra"); + config.setPassword("cassandra"); + config.setColumnFamily("observation"); + config.setKeyspace("airquality"); + + CassandraConnector connector = new CassandraConnector(config); + assertEquals("INSERT INTO airquality.observation (key, observed) VALUES (?, ?)", connector.getPreparedStatement().getQueryString()); + } + + @Test + public final void getReadingPreparedStatementTest() { + config = new CassandraSinkConfig(); + config.setRoots("localhost"); + config.setUserName("cassandra"); + config.setPassword("cassandra"); + config.setColumnFamily("reading"); + config.setKeyspace("airquality"); + + CassandraConnector connector = new CassandraConnector(config); + assertEquals("INSERT INTO airquality.reading " + + "(reporting_area, avg_ozone, avg_pm10, avg_pm25, date_observed, hour_observed, latitude, " + + "local_time_zone, longitude, max_ozone, max_pm10, max_pm25, min_ozone, min_pm10, min_pm25, " + + "readingid, state_code) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", + connector.getPreparedStatement().getQueryString()); + } +} diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java new file mode 100644 index 0000000000000..816e371434467 --- /dev/null +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java @@ -0,0 +1,98 @@ +/** + * 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.pulsar.io.cassandra; + +import java.io.File; +import java.io.FileInputStream; +import java.io.FileNotFoundException; +import java.util.Collections; +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.TimeUnit; + +import org.apache.pulsar.common.io.SinkConfig; +import org.apache.pulsar.functions.LocalRunner; +import org.apache.pulsar.io.cassandra.producers.InputTopicProducerThread; +import org.apache.pulsar.io.cassandra.producers.ReadingSchemaRecordProducer; +import org.yaml.snakeyaml.Yaml; + +/** + * Useful for testing within IDE. + * + */ +public class CassandraSinkExec { + + public static final String BROKER_URL = "pulsar://localhost:6650"; + // public static final String INPUT_TOPIC = "persistent://public/default/cassandra-observation-avro"; + // public static final String INPUT_TOPIC = "persistent://public/default/cassandra-observation-string"; + public static final String INPUT_TOPIC = "persistent://public/default/air-quality-reading-avro-3"; + public static final String CONFIG_FILE = "cassandra-sink-config.yaml"; + + public static void main(String[] args) throws Exception { + + SinkConfig config = getSinkConfig(); + + final LocalRunner localRunner = + LocalRunner.builder() + .brokerServiceUrl(BROKER_URL) + .sinkConfig(config) + .build(); + + localRunner.start(false); + + sendData(); + TimeUnit.MINUTES.sleep(10); + + localRunner.stop(); + + System.exit(0); + } + + private static SinkConfig getSinkConfig() throws FileNotFoundException { + SinkConfig sinkConfig = SinkConfig.builder() + .autoAck(true) + .cleanupSubscription(Boolean.TRUE) + .configs(getConfigs()) + .className(CassandraSink.class.getName()) + .inputs(Collections.singletonList(INPUT_TOPIC)) + .name("CassandraSink") + .build(); + + return sinkConfig; + } + + private static Map getConfigs() throws FileNotFoundException { + Map configs = new HashMap(); + + ClassLoader classLoader = Thread.currentThread().getContextClassLoader(); + File file = new File(classLoader.getResource(CONFIG_FILE).getFile()); + + Yaml yaml = new Yaml(); + configs = yaml.load(new FileInputStream(file)); + + return configs; + } + + private static final void sendData() throws InterruptedException { + TimeUnit.SECONDS.sleep(10); + InputTopicProducerThread writer = new ReadingSchemaRecordProducer(BROKER_URL, INPUT_TOPIC); + writer.run(); + } +} diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/AbstractGenericRecordProducer.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/AbstractGenericRecordProducer.java new file mode 100644 index 0000000000000..36c5c2d2fd955 --- /dev/null +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/AbstractGenericRecordProducer.java @@ -0,0 +1,43 @@ +/** + * 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.pulsar.io.cassandra.producers; + +import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.client.api.schema.*; +import org.apache.pulsar.common.schema.SchemaInfo; + + +public abstract class AbstractGenericRecordProducer extends InputTopicProducerThread { + + public AbstractGenericRecordProducer(String brokerUrl, String inputTopic) { + super(brokerUrl, inputTopic); + } + + @Override + Schema getSchema() { + return Schema.generic(getGenericSchemaInfo()); + } + + @Override + abstract GenericRecord getValue(); + + abstract SchemaInfo getGenericSchemaInfo(); + + } diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicProducerThread.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicProducerThread.java new file mode 100644 index 0000000000000..c43a24fb1248f --- /dev/null +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicProducerThread.java @@ -0,0 +1,83 @@ +/** + * 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.pulsar.io.cassandra.producers; + +import org.apache.pulsar.client.api.Producer; +import org.apache.pulsar.client.api.PulsarClient; +import org.apache.pulsar.client.api.PulsarClientException; +import org.apache.pulsar.client.api.Schema; + +import java.util.Random; + +public abstract class InputTopicProducerThread implements Runnable { + + private Random rnd = new Random(); + + final String inputTopic; + final String brokerUrl; + PulsarClient client; + Producer producer; + + public InputTopicProducerThread(String brokerUrl, String inputTopic) { + this.brokerUrl = brokerUrl; + this.inputTopic = inputTopic; + } + + @Override + public void run() { + for (int idx = 0; idx < 100; idx++) { + try { + getProducer().newMessage().key(getKey()).value(getValue()).send(); + } catch (PulsarClientException e) { + e.printStackTrace(); + } + } + } + + String getKey() { + Integer i = Integer.valueOf(rnd.nextInt(999999)); + return i.toString(); + } + + abstract T getValue(); + + abstract Schema getSchema(); + + private PulsarClient getPulsarClient() throws PulsarClientException { + if (client == null) { + client = PulsarClient.builder() + .serviceUrl(brokerUrl) + .build(); + } + + return client; + } + + private Producer getProducer() throws PulsarClientException { + if (producer == null) { + producer = getPulsarClient().newProducer(getSchema()) + .topic(inputTopic) + .create(); + } + + return producer; + } + +} diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java new file mode 100644 index 0000000000000..2e32fe746199b --- /dev/null +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java @@ -0,0 +1,43 @@ +/** + * 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.pulsar.io.cassandra.producers; + +import org.apache.pulsar.client.api.Schema; + +import java.util.Random; + +public class InputTopicStringProducer extends InputTopicProducerThread { + + public InputTopicStringProducer(String brokerUrl, String inputTopic) { + super(brokerUrl, inputTopic); + } + + @Override + String getValue() { + String val = "{\"dateObserved\":\"2022-06-08 \",\"hourObserved\":11,\"localTimeZone\":\"PST\",\"reportingArea\":\"Redwood City\",\"stateCode\":\"CA\",\"latitude\":37.48,\"longitude\":-122.22,\"parameterName\":\"PM2.5\",\"aqi\":20,\"category\":{\"number\":1,\"name\":\"Good\",\"additionalProperties\":{}}"; + return val; + } + + @Override + Schema getSchema() { + return Schema.STRING; + } + +} diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ObservationSchemaRecordProducer.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ObservationSchemaRecordProducer.java new file mode 100644 index 0000000000000..01caf13d80a12 --- /dev/null +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ObservationSchemaRecordProducer.java @@ -0,0 +1,60 @@ +/** + * 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.pulsar.io.cassandra.producers; + +import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.client.api.schema.GenericRecord; +import org.apache.pulsar.client.api.schema.RecordSchemaBuilder; +import org.apache.pulsar.client.api.schema.SchemaBuilder; +import org.apache.pulsar.common.schema.SchemaInfo; +import org.apache.pulsar.common.schema.SchemaType; + +public class ObservationSchemaRecordProducer extends AbstractGenericRecordProducer { + + public ObservationSchemaRecordProducer(String brokerUrl, String inputTopic) { + super(brokerUrl, inputTopic); + } + + @Override + GenericRecord getValue() { + String val = "Some random string"; + + GenericRecord record = Schema.generic(getGenericSchemaInfo()).newRecordBuilder() + .set("key", getKey()) + .set("observed", val) + .build(); + + return record; + } + + @Override + SchemaInfo getGenericSchemaInfo() { + RecordSchemaBuilder recordSchemaBuilder = + SchemaBuilder.record("airquality.observation"); + + recordSchemaBuilder.field("key") + .type(SchemaType.STRING).required(); + + recordSchemaBuilder.field("observed") + .type(SchemaType.STRING).optional(); + + return recordSchemaBuilder.build(SchemaType.AVRO); + } +} diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ReadingSchemaRecordProducer.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ReadingSchemaRecordProducer.java new file mode 100644 index 0000000000000..c5f8541228a83 --- /dev/null +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ReadingSchemaRecordProducer.java @@ -0,0 +1,93 @@ +/** + * 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.pulsar.io.cassandra.producers; + +import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.client.api.schema.GenericRecord; +import org.apache.pulsar.client.api.schema.RecordSchemaBuilder; +import org.apache.pulsar.client.api.schema.SchemaBuilder; +import org.apache.pulsar.common.schema.SchemaInfo; +import org.apache.pulsar.common.schema.SchemaType; + +import java.util.Random; + +public class ReadingSchemaRecordProducer extends AbstractGenericRecordProducer { + + private Random rnd = new Random(); + + int lastReadingId = rnd.nextInt(50000); + + public ReadingSchemaRecordProducer(String brokerUrl, String inputTopic) { + super(brokerUrl, inputTopic); + } + + @Override + GenericRecord getValue() { + GenericRecord record = Schema.generic(getGenericSchemaInfo()) + .newRecordBuilder() + .set("readingid", lastReadingId++ + "") + .set("avg_ozone", rnd.nextDouble()) + .set("min_ozone", rnd.nextDouble()) + .set("max_ozone", rnd.nextDouble()) + .set("avg_pm10", rnd.nextDouble()) + .set("min_pm10", rnd.nextDouble()) + .set("max_pm10", rnd.nextDouble()) + .set("avg_pm25", rnd.nextDouble()) + .set("min_pm25", rnd.nextDouble()) + .set("max_pm25", rnd.nextDouble()) + .set("local_time_zone", "PST") + .set("state_code", "CA") + .set("reporting_area", lastReadingId + "") + .set("hour_observed", rnd.nextInt(24)) + .set("date_observed", "2022-06-18") + .set("latitude", Float.valueOf(40.021f)) + .set("longitude", Float.valueOf(-122.33f)) + .build(); + + return record; + } + + @Override + SchemaInfo getGenericSchemaInfo() { + RecordSchemaBuilder schemaBuilder = + SchemaBuilder.record("airquality.reading"); + + schemaBuilder.field("readingid").type(SchemaType.STRING).required(); + schemaBuilder.field("avg_ozone").type(SchemaType.DOUBLE); + schemaBuilder.field("min_ozone").type(SchemaType.DOUBLE); + schemaBuilder.field("max_ozone").type(SchemaType.DOUBLE); + schemaBuilder.field("avg_pm10").type(SchemaType.DOUBLE); + schemaBuilder.field("min_pm10").type(SchemaType.DOUBLE); + schemaBuilder.field("max_pm10").type(SchemaType.DOUBLE); + schemaBuilder.field("avg_pm25").type(SchemaType.DOUBLE); + schemaBuilder.field("min_pm25").type(SchemaType.DOUBLE); + schemaBuilder.field("max_pm25").type(SchemaType.DOUBLE); + + schemaBuilder.field("local_time_zone").type(SchemaType.STRING); + schemaBuilder.field("state_code").type(SchemaType.STRING); + schemaBuilder.field("reporting_area").type(SchemaType.STRING).required(); + schemaBuilder.field("hour_observed").type(SchemaType.INT32); + schemaBuilder.field("date_observed").type(SchemaType.STRING); + schemaBuilder.field("latitude").type(SchemaType.FLOAT); + schemaBuilder.field("longitude").type(SchemaType.FLOAT); + + return schemaBuilder.build(SchemaType.AVRO); + } +} diff --git a/pulsar-io/cassandra/src/test/resources/cassandra-sink-config.yaml b/pulsar-io/cassandra/src/test/resources/cassandra-sink-config.yaml new file mode 100644 index 0000000000000..08ef19681c93e --- /dev/null +++ b/pulsar-io/cassandra/src/test/resources/cassandra-sink-config.yaml @@ -0,0 +1,23 @@ +# 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. + +roots: "localhost" +keyspace : "airquality" +columnFamily: "reading" +userName : "cassandra" +password : "cassandra" + diff --git a/tests/integration/pom.xml b/tests/integration/pom.xml index 477242799d4ce..3d946e9909c4c 100644 --- a/tests/integration/pom.xml +++ b/tests/integration/pom.xml @@ -81,7 +81,8 @@ com.datastax.cassandra cassandra-driver-core - + 3.11.2 + io.netty netty-handler From f13c640e532d0a07a2b45b0b1af53286b5bde518 Mon Sep 17 00:00:00 2001 From: david-streamlio Date: Mon, 20 Jun 2022 11:24:33 -0700 Subject: [PATCH 02/14] Added support for raw JSON String --- pulsar-io/cassandra/pom.xml | 12 +++ .../io/cassandra/CassandraAbstractSink.java | 90 ------------------- ...k.java => CassandraGenericRecordSink.java} | 25 ++---- .../pulsar/io/cassandra/CassandraSink.java | 56 ++++++------ .../io/cassandra/CassandraSinkConfig.java | 16 +--- .../{ => util}/CassandraConnector.java | 9 +- .../io/cassandra/CassandraSinkExec.java | 8 +- .../producers/InputTopicStringProducer.java | 35 +++++++- .../{ => util}/CassandraConnectorTest.java | 4 +- 9 files changed, 99 insertions(+), 156 deletions(-) delete mode 100644 pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraAbstractSink.java rename pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/{CassandraStringSink.java => CassandraGenericRecordSink.java} (54%) rename pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/{ => util}/CassandraConnector.java (93%) rename pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/{ => util}/CassandraConnectorTest.java (92%) diff --git a/pulsar-io/cassandra/pom.xml b/pulsar-io/cassandra/pom.xml index abf18b9c638bc..3f999448522d7 100644 --- a/pulsar-io/cassandra/pom.xml +++ b/pulsar-io/cassandra/pom.xml @@ -60,6 +60,12 @@ + + commons-beanutils + commons-beanutils + 1.9.4 + + org.apache.pulsar pulsar-functions-local-runner-original @@ -73,6 +79,12 @@ 2.11.0-SNAPSHOT compile + + commons-beanutils + commons-beanutils + 1.9.4 + compile + diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraAbstractSink.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraAbstractSink.java deleted file mode 100644 index 1323379c66c9b..0000000000000 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraAbstractSink.java +++ /dev/null @@ -1,90 +0,0 @@ -/** - * 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.pulsar.io.cassandra; - -import com.datastax.driver.core.*; -import com.datastax.driver.core.ResultSet; -import com.google.common.util.concurrent.FutureCallback; -import com.google.common.util.concurrent.Futures; -import com.google.common.util.concurrent.MoreExecutors; -import java.util.Map; -import org.apache.pulsar.functions.api.Record; -import org.apache.pulsar.io.common.IOConfigUtils; -import org.apache.pulsar.io.core.KeyValue; -import org.apache.pulsar.io.core.Sink; -import org.apache.pulsar.io.core.SinkContext; - -/** - * A Simple abstract class for Cassandra sink. - * Users need to implement extractKeyValue function to use this sink - */ -public abstract class CassandraAbstractSink implements Sink { - - private CassandraConnector connector; - private CassandraSinkConfig cassandraSinkConfig; - - @Override - public void open(Map config, SinkContext sinkContext) throws Exception { - - cassandraSinkConfig = IOConfigUtils.loadWithSecrets(config, CassandraSinkConfig.class, sinkContext); - - if (cassandraSinkConfig.getRoots() == null - || cassandraSinkConfig.getKeyspace() == null - || cassandraSinkConfig.getKeyname() == null - || cassandraSinkConfig.getColumnFamily() == null - || cassandraSinkConfig.getColumnName() == null) { - throw new IllegalArgumentException("Required property not set."); - } - - connector = new CassandraConnector(cassandraSinkConfig); - connector.connect(); - } - - @Override - public void close() throws Exception { - connector.close(); - } - - @Override - public void write(Record record) { - - KeyValue keyValue = extractKeyValue(record); - - BoundStatement bound = connector.getPreparedStatement() - .bind(keyValue.getKey(), keyValue.getValue()); - - ResultSetFuture future = connector.getSession().executeAsync(bound); - - Futures.addCallback(future, - new FutureCallback() { - @Override - public void onSuccess(ResultSet result) { - record.ack(); - } - - @Override - public void onFailure(Throwable t) { - record.fail(); - } - }, MoreExecutors.directExecutor()); - } - - public abstract KeyValue extractKeyValue(Record record); -} \ No newline at end of file diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraStringSink.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraGenericRecordSink.java similarity index 54% rename from pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraStringSink.java rename to pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraGenericRecordSink.java index 3ff3af732e8d1..3ecc2ee4d575a 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraStringSink.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraGenericRecordSink.java @@ -19,24 +19,15 @@ package org.apache.pulsar.io.cassandra; +import org.apache.pulsar.client.api.schema.GenericRecord; import org.apache.pulsar.functions.api.Record; -import org.apache.pulsar.io.core.KeyValue; -import org.apache.pulsar.io.core.annotations.Connector; -import org.apache.pulsar.io.core.annotations.IOType; +import org.apache.pulsar.io.cassandra.util.GenericRecordWrapper; +import org.apache.pulsar.io.cassandra.util.RecordWrapper; + +public class CassandraGenericRecordSink extends CassandraSink { -/** - * Cassandra sink that treats incoming messages on the input topic as Strings - * and writes key/value pairs. - */ -@Connector( - name = "cassandra", - type = IOType.SINK, - help = "The CassandraStringSink is used for moving messages from Pulsar to Cassandra.", - configClass = CassandraSinkConfig.class) -public class CassandraStringSink extends CassandraAbstractSink { @Override - public KeyValue extractKeyValue(Record record) { - String key = record.getKey().orElseGet(() -> new String(record.getValue())); - return new KeyValue<>(key, new String(record.getValue())); + RecordWrapper wrapRecord(Record record) { + return new GenericRecordWrapper(record.getValue()); } -} \ No newline at end of file +} diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSink.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSink.java index 470d5a0267715..7f68fd9cad164 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSink.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSink.java @@ -19,24 +19,34 @@ package org.apache.pulsar.io.cassandra; -import com.datastax.driver.core.*; +import com.datastax.driver.core.BoundStatement; +import com.datastax.driver.core.PreparedStatement; +import com.datastax.driver.core.ResultSet; +import com.datastax.driver.core.ResultSetFuture; import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.MoreExecutors; -import org.apache.pulsar.client.api.schema.GenericRecord; import org.apache.pulsar.functions.api.Record; +import org.apache.pulsar.io.cassandra.util.*; import org.apache.pulsar.io.common.IOConfigUtils; import org.apache.pulsar.io.core.Sink; import org.apache.pulsar.io.core.SinkContext; +import org.apache.pulsar.io.core.annotations.Connector; +import org.apache.pulsar.io.core.annotations.IOType; import java.util.Map; -public class CassandraSink implements Sink { +@Connector( + name = "cassandra", + type = IOType.SINK, + help = "The CassandraStringSink is used for moving messages from Pulsar to Cassandra.", + configClass = CassandraSinkConfig.class) +public abstract class CassandraSink implements Sink { - private CassandraConnector connector; - private CassandraSinkConfig cassandraSinkConfig; - private PreparedStatement stmt; - private BoundStatement boundStatement; + CassandraConnector connector; + CassandraSinkConfig cassandraSinkConfig; + PreparedStatement stmt; + BoundStatementProvider boundStatementProvider; @Override public void open(Map config, SinkContext ctx) throws Exception { @@ -51,20 +61,19 @@ public void open(Map config, SinkContext ctx) throws Exception { connector = new CassandraConnector(cassandraSinkConfig); connector.connect(); + + boundStatementProvider = new BoundStatementProvider( + TableMetadataProvider.getTableDefinition( + connector.getTableMetadata(), + cassandraSinkConfig.getKeyspace(), + cassandraSinkConfig.getColumnFamily())); } @Override - public void write(Record record) throws Exception { - - Object[] boundValues = new Object[connector.getTableFields().size()]; - GenericRecord generic = record.getValue(); - - for (int idx = 0; idx < connector.getTableFields().size(); idx++) { - String fieldName = connector.getTableFields().get(idx); - boundValues[idx] = generic.getField(fieldName); - } + public void write(Record record) throws Exception { - BoundStatement bs = getBoundStatement().bind(boundValues); + BoundStatement bs = boundStatementProvider.bindStatement( + getStatement(), wrapRecord(record)); ResultSetFuture future = connector.getSession().executeAsync(bs); @@ -87,18 +96,13 @@ public void close() throws Exception { connector.close(); } - private PreparedStatement getStatement() { + abstract RecordWrapper wrapRecord(Record record); + + PreparedStatement getStatement() { if (stmt == null) { - stmt = connector.getPreparedStatement(); + stmt = connector.getPreparedStatement(); } return stmt; } - private BoundStatement getBoundStatement() { - if (boundStatement == null) { - boundStatement = getStatement().bind(); - boundStatement.setConsistencyLevel(ConsistencyLevel.ALL); - } - return boundStatement; - } } diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSinkConfig.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSinkConfig.java index c9b52e225284a..6c5ab941b9e44 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSinkConfig.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSinkConfig.java @@ -39,7 +39,7 @@ public class CassandraSinkConfig implements Serializable { required = false, defaultValue = "", sensitive = true, - help = "Username used to connect to the database specified by `jdbcUrl`" + help = "Username used to connect to the database specified by `root`" ) private String userName; @@ -47,7 +47,7 @@ public class CassandraSinkConfig implements Serializable { required = false, defaultValue = "", sensitive = true, - help = "Password used to connect to the database specified by `jdbcUrl`" + help = "Password used to connect to the database specified by `root`" ) private String password; @@ -63,24 +63,12 @@ public class CassandraSinkConfig implements Serializable { help = "The key space used for writing pulsar messages to") private String keyspace; - @FieldDoc( - required = true, - defaultValue = "", - help = "The key name of the cassandra column family") - private String keyname; - @FieldDoc( required = true, defaultValue = "", help = "The cassandra column family name") private String columnFamily; - @FieldDoc( - required = true, - defaultValue = "", - help = "The column name of the cassandra column family") - private String columnName; - public static CassandraSinkConfig load(String yamlFile) throws IOException { ObjectMapper mapper = new ObjectMapper(new YAMLFactory()); return mapper.readValue(new File(yamlFile), CassandraSinkConfig.class); diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraConnector.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/CassandraConnector.java similarity index 93% rename from pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraConnector.java rename to pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/CassandraConnector.java index b4dede9eaf053..e83fb39d80bfd 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraConnector.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/CassandraConnector.java @@ -17,10 +17,11 @@ * under the License. */ -package org.apache.pulsar.io.cassandra; +package org.apache.pulsar.io.cassandra.util; import com.datastax.driver.core.*; import com.datastax.driver.core.Session; +import org.apache.pulsar.io.cassandra.CassandraSinkConfig; import java.util.ArrayList; import java.util.List; @@ -42,12 +43,16 @@ public void connect() { } public Session getSession() { - if (session == null) { + if (session == null || session.isClosed()) { this.connect(); } return session; } + public Metadata getTableMetadata() { + return getCluster().getMetadata(); + } + public PreparedStatement getPreparedStatement() { if (statement == null) { List fields = getTableFields(); diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java index 816e371434467..45ac58968cb42 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java @@ -30,6 +30,7 @@ import org.apache.pulsar.common.io.SinkConfig; import org.apache.pulsar.functions.LocalRunner; import org.apache.pulsar.io.cassandra.producers.InputTopicProducerThread; +import org.apache.pulsar.io.cassandra.producers.InputTopicStringProducer; import org.apache.pulsar.io.cassandra.producers.ReadingSchemaRecordProducer; import org.yaml.snakeyaml.Yaml; @@ -40,9 +41,8 @@ public class CassandraSinkExec { public static final String BROKER_URL = "pulsar://localhost:6650"; - // public static final String INPUT_TOPIC = "persistent://public/default/cassandra-observation-avro"; - // public static final String INPUT_TOPIC = "persistent://public/default/cassandra-observation-string"; - public static final String INPUT_TOPIC = "persistent://public/default/air-quality-reading-avro-3"; + public static final String INPUT_TOPIC = "persistent://public/default/air-quality-reading-generic"; + // public static final String INPUT_TOPIC = "persistent://public/default/air-quality-reading-string"; public static final String CONFIG_FILE = "cassandra-sink-config.yaml"; public static void main(String[] args) throws Exception { @@ -70,7 +70,7 @@ private static SinkConfig getSinkConfig() throws FileNotFoundException { .autoAck(true) .cleanupSubscription(Boolean.TRUE) .configs(getConfigs()) - .className(CassandraSink.class.getName()) + .className(CassandraGenericRecordSink.class.getName()) .inputs(Collections.singletonList(INPUT_TOPIC)) .name("CassandraSink") .build(); diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java index 2e32fe746199b..1e867bd2eef67 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java @@ -19,20 +19,51 @@ package org.apache.pulsar.io.cassandra.producers; +import com.google.gson.Gson; +import com.google.gson.reflect.TypeToken; import org.apache.pulsar.client.api.Schema; +import java.lang.reflect.Type; +import java.util.HashMap; import java.util.Random; +import java.util.SortedMap; +import java.util.TreeMap; public class InputTopicStringProducer extends InputTopicProducerThread { + private Random rnd = new Random(); + int lastReadingId = rnd.nextInt(900000); + public InputTopicStringProducer(String brokerUrl, String inputTopic) { super(brokerUrl, inputTopic); } @Override String getValue() { - String val = "{\"dateObserved\":\"2022-06-08 \",\"hourObserved\":11,\"localTimeZone\":\"PST\",\"reportingArea\":\"Redwood City\",\"stateCode\":\"CA\",\"latitude\":37.48,\"longitude\":-122.22,\"parameterName\":\"PM2.5\",\"aqi\":20,\"category\":{\"number\":1,\"name\":\"Good\",\"additionalProperties\":{}}"; - return val; + + SortedMap elements = new TreeMap(); + elements.put("readingid", lastReadingId++ + ""); + elements.put("avg_ozone", rnd.nextDouble()); + elements.put("min_ozone", rnd.nextDouble()); + elements.put("max_ozone", rnd.nextDouble()); + elements.put("avg_pm10", rnd.nextDouble()); + elements.put("min_pm10", rnd.nextDouble()); + elements.put("max_pm10", rnd.nextDouble()); + elements.put("avg_pm25", rnd.nextDouble()); + elements.put("min_pm25", rnd.nextDouble()); + elements.put("max_pm25", rnd.nextDouble()); + elements.put("local_time_zone", "PST"); + elements.put("state_code", "CA"); + elements.put("reporting_area", lastReadingId + ""); + elements.put("hour_observed", rnd.nextInt(24)); + elements.put("date_observed", "2022-06-18"); + elements.put("latitude", Float.valueOf(40.021f)); + elements.put("longitude", Float.valueOf(-122.33f)); + + Gson gson = new Gson(); + Type gsonType = new TypeToken(){}.getType(); + String gsonString = gson.toJson(elements,gsonType); + return gsonString; } @Override diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraConnectorTest.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java similarity index 92% rename from pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraConnectorTest.java rename to pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java index 91c063a9ee2e0..45cecc6797078 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraConnectorTest.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java @@ -1,5 +1,7 @@ -package org.apache.pulsar.io.cassandra; +package org.apache.pulsar.io.cassandra.util; +import org.apache.pulsar.io.cassandra.CassandraSinkConfig; +import org.apache.pulsar.io.cassandra.util.CassandraConnector; import org.junit.Ignore; import org.junit.Test; From fb7beb70b3e88e8ccaa64b50f1ba77aa563b3818 Mon Sep 17 00:00:00 2001 From: david-streamlio Date: Wed, 22 Jun 2022 07:39:54 -0700 Subject: [PATCH 03/14] Added RecordWrapper interface --- .../io/cassandra/CassandraStringSink.java | 32 ++++++ .../util/BoundStatementProvider.java | 50 +++++++++ .../cassandra/util/GenericRecordWrapper.java | 39 +++++++ .../io/cassandra/util/RecordWrapper.java | 49 +++++++++ .../cassandra/util/StringRecordWrapper.java | 47 ++++++++ .../cassandra/util/TableMetadataProvider.java | 102 ++++++++++++++++++ .../util/TableMetadataProviderTest.java | 54 ++++++++++ 7 files changed, 373 insertions(+) create mode 100644 pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraStringSink.java create mode 100644 pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/BoundStatementProvider.java create mode 100644 pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/GenericRecordWrapper.java create mode 100644 pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/RecordWrapper.java create mode 100644 pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/StringRecordWrapper.java create mode 100644 pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/TableMetadataProvider.java create mode 100644 pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraStringSink.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraStringSink.java new file mode 100644 index 0000000000000..2cbd9647701cf --- /dev/null +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraStringSink.java @@ -0,0 +1,32 @@ +/** + * 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.pulsar.io.cassandra; + +import org.apache.pulsar.functions.api.Record; +import org.apache.pulsar.io.cassandra.util.RecordWrapper; +import org.apache.pulsar.io.cassandra.util.StringRecordWrapper; + +public class CassandraStringSink extends CassandraSink { + + @Override + RecordWrapper wrapRecord(Record record) { + return new StringRecordWrapper(record.getValue()); + } +} diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/BoundStatementProvider.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/BoundStatementProvider.java new file mode 100644 index 0000000000000..8c73eaf7a16bd --- /dev/null +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/BoundStatementProvider.java @@ -0,0 +1,50 @@ +/** + * 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.pulsar.io.cassandra.util; + +import com.datastax.driver.core.BoundStatement; +import com.datastax.driver.core.PreparedStatement; +import org.apache.commons.beanutils.converters.*; + + +public class BoundStatementProvider { + + NumberConverter converter = new IntegerConverter(); + + final TableMetadataProvider.TableDefinition tableDefinition; + + public BoundStatementProvider(TableMetadataProvider.TableDefinition tableDefinition) { + this.tableDefinition = tableDefinition; + } + + public BoundStatement bindStatement(PreparedStatement stmt, RecordWrapper wrapper) { + Object[] boundValues = new Object[tableDefinition.getColumns().size()]; + int idx = 0; + + for (TableMetadataProvider.ColumnId column : tableDefinition.getColumns()) { + if (wrapper.containsKey(column.getName())) { + boundValues[idx] = wrapper.get(column); + } + idx++; + } + return stmt.bind(boundValues); + } + +} diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/GenericRecordWrapper.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/GenericRecordWrapper.java new file mode 100644 index 0000000000000..3e1b24b16ff20 --- /dev/null +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/GenericRecordWrapper.java @@ -0,0 +1,39 @@ +/** + * 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.pulsar.io.cassandra.util; + +import org.apache.pulsar.client.api.schema.GenericRecord; + +public class GenericRecordWrapper extends RecordWrapper { + + public GenericRecordWrapper(GenericRecord value) { + super(value); + } + + @Override + public Object get(TableMetadataProvider.ColumnId column) { + return getValueAsExpectedType(recordValue.getField(column.getName()), column); + } + + @Override + public boolean containsKey(String name) { + return this.recordValue.getField(name) != null; + } +} diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/RecordWrapper.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/RecordWrapper.java new file mode 100644 index 0000000000000..b6a1b2b5bc429 --- /dev/null +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/RecordWrapper.java @@ -0,0 +1,49 @@ +/** + * 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.pulsar.io.cassandra.util; + +import org.apache.commons.beanutils.converters.IntegerConverter; +import org.apache.commons.beanutils.converters.NumberConverter; + +public abstract class RecordWrapper { + + T recordValue; + NumberConverter converter = new IntegerConverter(); + + public RecordWrapper(T value) { + this.recordValue = value; + } + + public abstract Object get(TableMetadataProvider.ColumnId column); + + public abstract boolean containsKey(String name); + + Object getValueAsExpectedType(Object value, TableMetadataProvider.ColumnId column) { + switch (column.getType().getName()) { + case FLOAT: return converter.convert(Float.class, value); + case INT: return converter.convert(Integer.class, value); + case DOUBLE: return converter.convert(Double.class, value); + case TEXT: return value.toString(); + default: return value; + } + + } + +} diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/StringRecordWrapper.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/StringRecordWrapper.java new file mode 100644 index 0000000000000..febaba4832eea --- /dev/null +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/StringRecordWrapper.java @@ -0,0 +1,47 @@ +/** + * 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.pulsar.io.cassandra.util; + +import com.google.gson.Gson; +import com.google.gson.reflect.TypeToken; + +import java.util.HashMap; +import java.util.Map; + +public class StringRecordWrapper extends RecordWrapper { + + private Map valuesMap; + + public StringRecordWrapper(String jsonString) { + super(jsonString); + valuesMap = new Gson().fromJson(jsonString, + new TypeToken>() {}.getType()); + } + + @Override + public Object get(TableMetadataProvider.ColumnId column) { + return getValueAsExpectedType(valuesMap.get(column.getName()), column); + } + + @Override + public boolean containsKey(String name) { + return valuesMap.containsKey(name); + } +} diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/TableMetadataProvider.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/TableMetadataProvider.java new file mode 100644 index 0000000000000..c3e2cec65f194 --- /dev/null +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/TableMetadataProvider.java @@ -0,0 +1,102 @@ +/** + * 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.pulsar.io.cassandra.util; + +import com.datastax.driver.core.*; +import com.google.common.collect.Lists; +import lombok.*; + +import java.util.List; + +public class TableMetadataProvider { + + @Data(staticConstructor = "of") + public static class TableId { + private final String keyspaceName; + private final String columnFamily; + } + + @Data(staticConstructor = "of") + public static class ColumnId { + private final TableId tableId; + private final String name; + private final DataType type; + } + + @Setter + @Getter + @EqualsAndHashCode + @ToString + public static class TableDefinition { + private final TableId tableId; + private final List columns; + private final List partitionKeyColumns; + private final List primaryKeyColumns; + + private TableDefinition(TableId tableId, List columns) { + this(tableId, columns, null, null); + } + private TableDefinition(TableId tableId, List columns, + List partitionKeyColumns, List primaryKeyColumns) { + this.tableId = tableId; + this.columns = columns; + this.partitionKeyColumns = partitionKeyColumns; + this.primaryKeyColumns = primaryKeyColumns; + } + + public static TableDefinition of(TableId tableId, List columns) { + return new TableDefinition(tableId, columns); + } + + public static TableDefinition of(TableId tableId, List columns, + List nonKeyColumns, List keyColumns) { + return new TableDefinition(tableId, columns, nonKeyColumns, keyColumns); + } + + } + + public static TableDefinition getTableDefinition(Metadata clusterMetadata, String keyspace, String columnFamily) { + + TableId tableId = TableId.of(keyspace, columnFamily); + TableDefinition table = TableDefinition.of(tableId, + Lists.newArrayList(), Lists.newArrayList(), Lists.newArrayList()); + + TableMetadata meta = clusterMetadata + .getKeyspace(keyspace) + .getTable(columnFamily); + + for (ColumnMetadata col : meta.getPrimaryKey()) { + ColumnId columnId = ColumnId.of(tableId, col.getName(), col.getType()); + table.getPrimaryKeyColumns().add(columnId); + } + + for (ColumnMetadata col : meta.getPartitionKey()) { + ColumnId columnId = ColumnId.of(tableId, col.getName(), col.getType()); + table.getPartitionKeyColumns().add(columnId); + } + + for (ColumnMetadata col : meta.getColumns()) { + ColumnId columnId = ColumnId.of(tableId, col.getName(), col.getType()); + table.getColumns().add(columnId); + } + + return table; + } +} diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java new file mode 100644 index 0000000000000..1c6f8ca42b98a --- /dev/null +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java @@ -0,0 +1,54 @@ +/** + * 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.pulsar.io.cassandra.util; + +import org.apache.pulsar.io.cassandra.CassandraSinkConfig; +import org.testng.annotations.Ignore; +import org.testng.annotations.Test; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; + +public class TableMetadataProviderTest { + + @Test + @Ignore + public final void getTableDefinitionTest() { + + CassandraSinkConfig config = new CassandraSinkConfig(); + config.setRoots("localhost"); + config.setUserName("cassandra"); + config.setPassword("cassandra"); + config.setColumnFamily("observation"); + config.setKeyspace("airquality"); + + CassandraConnector connector = new CassandraConnector(config); + + TableMetadataProvider.TableDefinition table = + TableMetadataProvider.getTableDefinition( + connector.getTableMetadata(), + "airquality", "reading"); + + assertNotNull(table); + assertEquals(17, table.getColumns().size()); + assertEquals(1, table.getPartitionKeyColumns().size()); + assertEquals(1, table.getPrimaryKeyColumns().size()); + } +} From 608f32f053bfdd05d2657d0e6d406ea6ef3875fa Mon Sep 17 00:00:00 2001 From: david-streamlio <35466513+david-streamlio@users.noreply.github.com> Date: Thu, 13 Oct 2022 08:54:19 -0700 Subject: [PATCH 04/14] fixed issue from code review --- pulsar-io/cassandra/pom.xml | 11 +++------- .../pulsar/io/cassandra/CassandraSink.java | 18 +++++++++++----- .../util/BoundStatementProvider.java | 5 +++-- .../io/cassandra/util/CassandraConnector.java | 21 +++++++++++-------- .../cassandra/util/StringRecordWrapper.java | 4 ++-- .../cassandra/util/TableMetadataProvider.java | 12 ++++++++--- .../io/cassandra/util/package-info.java | 19 +++++++++++++++++ .../io/cassandra/CassandraSinkExec.java | 17 ++++++++------- .../AbstractGenericRecordProducer.java | 2 +- .../producers/InputTopicProducerThread.java | 8 ++++--- .../producers/InputTopicStringProducer.java | 1 + .../ObservationSchemaRecordProducer.java | 1 + .../ReadingSchemaRecordProducer.java | 1 + .../util/CassandraConnectorTest.java | 20 ++++++++++++++++++ .../util/TableMetadataProviderTest.java | 4 ++-- .../test/resources/cassandra-sink-config.yaml | 2 ++ 16 files changed, 104 insertions(+), 42 deletions(-) create mode 100644 pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/package-info.java diff --git a/pulsar-io/cassandra/pom.xml b/pulsar-io/cassandra/pom.xml index 3f999448522d7..723919e34b97b 100644 --- a/pulsar-io/cassandra/pom.xml +++ b/pulsar-io/cassandra/pom.xml @@ -60,25 +60,20 @@ - - commons-beanutils - commons-beanutils - 1.9.4 - - org.apache.pulsar pulsar-functions-local-runner-original - 2.11.0-SNAPSHOT + ${project.version} test org.apache.pulsar pulsar-io-common - 2.11.0-SNAPSHOT + ${project.version} compile + commons-beanutils commons-beanutils diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSink.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSink.java index 7f68fd9cad164..d498c5751b9e1 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSink.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSink.java @@ -26,16 +26,18 @@ import com.google.common.util.concurrent.FutureCallback; import com.google.common.util.concurrent.Futures; import com.google.common.util.concurrent.MoreExecutors; +import java.util.Map; import org.apache.pulsar.functions.api.Record; -import org.apache.pulsar.io.cassandra.util.*; +import org.apache.pulsar.io.cassandra.util.BoundStatementProvider; +import org.apache.pulsar.io.cassandra.util.CassandraConnector; +import org.apache.pulsar.io.cassandra.util.RecordWrapper; +import org.apache.pulsar.io.cassandra.util.TableMetadataProvider; import org.apache.pulsar.io.common.IOConfigUtils; import org.apache.pulsar.io.core.Sink; import org.apache.pulsar.io.core.SinkContext; import org.apache.pulsar.io.core.annotations.Connector; import org.apache.pulsar.io.core.annotations.IOType; -import java.util.Map; - @Connector( name = "cassandra", type = IOType.SINK, @@ -92,8 +94,14 @@ public void onFailure(Throwable t) { } @Override - public void close() throws Exception { - connector.close(); + public void close() { + if (connector != null) { + try { + connector.close(); + } catch (final Throwable t) { + + } + } } abstract RecordWrapper wrapRecord(Record record); diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/BoundStatementProvider.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/BoundStatementProvider.java index 8c73eaf7a16bd..1b0b03ba0dbe1 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/BoundStatementProvider.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/BoundStatementProvider.java @@ -21,9 +21,10 @@ import com.datastax.driver.core.BoundStatement; import com.datastax.driver.core.PreparedStatement; -import org.apache.commons.beanutils.converters.*; - +import org.apache.commons.beanutils.converters.IntegerConverter; +import org.apache.commons.beanutils.converters.NumberConverter; +@SuppressWarnings("rawtypes") public class BoundStatementProvider { NumberConverter converter = new IntegerConverter(); diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/CassandraConnector.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/CassandraConnector.java index e83fb39d80bfd..20b8cfd200194 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/CassandraConnector.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/CassandraConnector.java @@ -19,12 +19,15 @@ package org.apache.pulsar.io.cassandra.util; -import com.datastax.driver.core.*; +import com.datastax.driver.core.Cluster; +import com.datastax.driver.core.ColumnMetadata; +import com.datastax.driver.core.Metadata; +import com.datastax.driver.core.PreparedStatement; import com.datastax.driver.core.Session; -import org.apache.pulsar.io.cassandra.CassandraSinkConfig; - +import com.datastax.driver.core.TableMetadata; import java.util.ArrayList; import java.util.List; +import org.apache.pulsar.io.cassandra.CassandraSinkConfig; public class CassandraConnector { @@ -42,7 +45,7 @@ public void connect() { session = getCluster().connect(config.getKeyspace()); } - public Session getSession() { + public synchronized Session getSession() { if (session == null || session.isClosed()) { this.connect(); } @@ -62,7 +65,7 @@ public PreparedStatement getPreparedStatement() { for (int idx = 0; idx < fields.size(); idx++) { sb.append(fields.get(idx)); - if (idx < fields.size()-1) { + if (idx < fields.size() - 1) { sb.append(", "); } } @@ -71,7 +74,7 @@ public PreparedStatement getPreparedStatement() { for (int idx = 0; idx < fields.size(); idx++) { sb.append("?"); - if (idx < fields.size()-1) { + if (idx < fields.size() - 1) { sb.append(", "); } } @@ -100,7 +103,7 @@ List getTableFields() { return tableFields; } - private Cluster getCluster() { + private synchronized Cluster getCluster() { if (cluster == null) { String[] hosts = config.getRoots().split(","); @@ -129,7 +132,7 @@ private Cluster getCluster() { } public void close() { - session.close(); - cluster.close(); + getSession().close(); + getCluster().close(); } } diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/StringRecordWrapper.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/StringRecordWrapper.java index febaba4832eea..260d752ba2f1c 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/StringRecordWrapper.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/StringRecordWrapper.java @@ -21,17 +21,17 @@ import com.google.gson.Gson; import com.google.gson.reflect.TypeToken; - import java.util.HashMap; import java.util.Map; public class StringRecordWrapper extends RecordWrapper { + private static final Gson gson = new Gson(); private Map valuesMap; public StringRecordWrapper(String jsonString) { super(jsonString); - valuesMap = new Gson().fromJson(jsonString, + valuesMap = gson.fromJson(jsonString, new TypeToken>() {}.getType()); } diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/TableMetadataProvider.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/TableMetadataProvider.java index c3e2cec65f194..955202a060961 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/TableMetadataProvider.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/TableMetadataProvider.java @@ -19,11 +19,17 @@ package org.apache.pulsar.io.cassandra.util; -import com.datastax.driver.core.*; +import com.datastax.driver.core.ColumnMetadata; +import com.datastax.driver.core.DataType; +import com.datastax.driver.core.Metadata; +import com.datastax.driver.core.TableMetadata; import com.google.common.collect.Lists; -import lombok.*; - import java.util.List; +import lombok.Data; +import lombok.EqualsAndHashCode; +import lombok.Getter; +import lombok.Setter; +import lombok.ToString; public class TableMetadataProvider { diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/package-info.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/package-info.java new file mode 100644 index 0000000000000..61699fd875df2 --- /dev/null +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/package-info.java @@ -0,0 +1,19 @@ +/** + * 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.pulsar.io.cassandra.util; \ No newline at end of file diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java index 45ac58968cb42..ee129e74107b7 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java @@ -21,7 +21,7 @@ import java.io.File; import java.io.FileInputStream; -import java.io.FileNotFoundException; +import java.io.IOException; import java.util.Collections; import java.util.HashMap; import java.util.Map; @@ -30,7 +30,6 @@ import org.apache.pulsar.common.io.SinkConfig; import org.apache.pulsar.functions.LocalRunner; import org.apache.pulsar.io.cassandra.producers.InputTopicProducerThread; -import org.apache.pulsar.io.cassandra.producers.InputTopicStringProducer; import org.apache.pulsar.io.cassandra.producers.ReadingSchemaRecordProducer; import org.yaml.snakeyaml.Yaml; @@ -38,11 +37,12 @@ * Useful for testing within IDE. * */ +@SuppressWarnings({"unchecked", "rawtypes"}) public class CassandraSinkExec { public static final String BROKER_URL = "pulsar://localhost:6650"; public static final String INPUT_TOPIC = "persistent://public/default/air-quality-reading-generic"; - // public static final String INPUT_TOPIC = "persistent://public/default/air-quality-reading-string"; + public static final String CONFIG_FILE = "cassandra-sink-config.yaml"; public static void main(String[] args) throws Exception { @@ -65,7 +65,7 @@ public static void main(String[] args) throws Exception { System.exit(0); } - private static SinkConfig getSinkConfig() throws FileNotFoundException { + private static SinkConfig getSinkConfig() throws IOException { SinkConfig sinkConfig = SinkConfig.builder() .autoAck(true) .cleanupSubscription(Boolean.TRUE) @@ -78,14 +78,17 @@ private static SinkConfig getSinkConfig() throws FileNotFoundException { return sinkConfig; } - private static Map getConfigs() throws FileNotFoundException { + private static Map getConfigs() throws IOException { Map configs = new HashMap(); ClassLoader classLoader = Thread.currentThread().getContextClassLoader(); File file = new File(classLoader.getResource(CONFIG_FILE).getFile()); - Yaml yaml = new Yaml(); - configs = yaml.load(new FileInputStream(file)); + try (FileInputStream fis = new FileInputStream(file)) { + configs = new Yaml().load(fis); + } catch (IOException ex) { + throw ex; + } return configs; } diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/AbstractGenericRecordProducer.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/AbstractGenericRecordProducer.java index 36c5c2d2fd955..2494756844b1f 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/AbstractGenericRecordProducer.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/AbstractGenericRecordProducer.java @@ -23,7 +23,7 @@ import org.apache.pulsar.client.api.schema.*; import org.apache.pulsar.common.schema.SchemaInfo; - +@SuppressWarnings({"unchecked", "rawtypes"}) public abstract class AbstractGenericRecordProducer extends InputTopicProducerThread { public AbstractGenericRecordProducer(String brokerUrl, String inputTopic) { diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicProducerThread.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicProducerThread.java index c43a24fb1248f..a80f4d2cf326a 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicProducerThread.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicProducerThread.java @@ -19,6 +19,7 @@ package org.apache.pulsar.io.cassandra.producers; +import lombok.extern.slf4j.Slf4j; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; @@ -26,10 +27,11 @@ import java.util.Random; +@SuppressWarnings({"unchecked", "rawtypes"}) +@Slf4j public abstract class InputTopicProducerThread implements Runnable { private Random rnd = new Random(); - final String inputTopic; final String brokerUrl; PulsarClient client; @@ -46,7 +48,7 @@ public void run() { try { getProducer().newMessage().key(getKey()).value(getValue()).send(); } catch (PulsarClientException e) { - e.printStackTrace(); + log.error("Unable to connect to Pulsar", e); } } } @@ -58,7 +60,7 @@ String getKey() { abstract T getValue(); - abstract Schema getSchema(); + abstract Schema getSchema(); private PulsarClient getPulsarClient() throws PulsarClientException { if (client == null) { diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java index 1e867bd2eef67..2282de98249e7 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java @@ -29,6 +29,7 @@ import java.util.SortedMap; import java.util.TreeMap; +@SuppressWarnings({"unchecked", "rawtypes"}) public class InputTopicStringProducer extends InputTopicProducerThread { private Random rnd = new Random(); diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ObservationSchemaRecordProducer.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ObservationSchemaRecordProducer.java index 01caf13d80a12..60cf753a42a1e 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ObservationSchemaRecordProducer.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ObservationSchemaRecordProducer.java @@ -26,6 +26,7 @@ import org.apache.pulsar.common.schema.SchemaInfo; import org.apache.pulsar.common.schema.SchemaType; +@SuppressWarnings({"unchecked", "rawtypes"}) public class ObservationSchemaRecordProducer extends AbstractGenericRecordProducer { public ObservationSchemaRecordProducer(String brokerUrl, String inputTopic) { diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ReadingSchemaRecordProducer.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ReadingSchemaRecordProducer.java index c5f8541228a83..126c8312de61b 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ReadingSchemaRecordProducer.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ReadingSchemaRecordProducer.java @@ -28,6 +28,7 @@ import java.util.Random; +@SuppressWarnings({"unchecked", "rawtypes"}) public class ReadingSchemaRecordProducer extends AbstractGenericRecordProducer { private Random rnd = new Random(); diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java index 45cecc6797078..6dde31dd67ad7 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java @@ -1,3 +1,21 @@ +/** + * 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.pulsar.io.cassandra.util; import org.apache.pulsar.io.cassandra.CassandraSinkConfig; @@ -26,6 +44,7 @@ public final void securedTest() { } @Test + @Ignore public final void getObservationPreparedStatementTest() { config = new CassandraSinkConfig(); config.setRoots("localhost"); @@ -39,6 +58,7 @@ public final void getObservationPreparedStatementTest() { } @Test + @Ignore public final void getReadingPreparedStatementTest() { config = new CassandraSinkConfig(); config.setRoots("localhost"); diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java index 1c6f8ca42b98a..f14f78473eeff 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java @@ -20,8 +20,8 @@ package org.apache.pulsar.io.cassandra.util; import org.apache.pulsar.io.cassandra.CassandraSinkConfig; -import org.testng.annotations.Ignore; -import org.testng.annotations.Test; +import org.junit.Ignore; +import org.junit.Test; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; diff --git a/pulsar-io/cassandra/src/test/resources/cassandra-sink-config.yaml b/pulsar-io/cassandra/src/test/resources/cassandra-sink-config.yaml index 08ef19681c93e..0ba3bdb878728 100644 --- a/pulsar-io/cassandra/src/test/resources/cassandra-sink-config.yaml +++ b/pulsar-io/cassandra/src/test/resources/cassandra-sink-config.yaml @@ -1,3 +1,4 @@ +# # 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 @@ -14,6 +15,7 @@ # KIND, either express or implied. See the License for the # specific language governing permissions and limitations # under the License. +# roots: "localhost" keyspace : "airquality" From c748715703102d2bf07494ba0c676217efdb8722 Mon Sep 17 00:00:00 2001 From: tison Date: Sun, 11 Dec 2022 09:45:18 +0800 Subject: [PATCH 05/14] license header Signed-off-by: tison --- .../apache/pulsar/io/cassandra/CassandraGenericRecordSink.java | 3 +-- .../pulsar/io/cassandra/util/BoundStatementProvider.java | 3 +-- .../apache/pulsar/io/cassandra/util/CassandraConnector.java | 3 +-- .../apache/pulsar/io/cassandra/util/GenericRecordWrapper.java | 3 +-- .../org/apache/pulsar/io/cassandra/util/RecordWrapper.java | 3 +-- .../apache/pulsar/io/cassandra/util/StringRecordWrapper.java | 3 +-- .../apache/pulsar/io/cassandra/util/TableMetadataProvider.java | 3 +-- .../java/org/apache/pulsar/io/cassandra/util/package-info.java | 2 +- .../java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java | 3 +-- .../io/cassandra/producers/AbstractGenericRecordProducer.java | 3 +-- .../io/cassandra/producers/InputTopicProducerThread.java | 3 +-- .../io/cassandra/producers/InputTopicStringProducer.java | 3 +-- .../cassandra/producers/ObservationSchemaRecordProducer.java | 3 +-- .../io/cassandra/producers/ReadingSchemaRecordProducer.java | 3 +-- .../pulsar/io/cassandra/util/CassandraConnectorTest.java | 2 +- .../pulsar/io/cassandra/util/TableMetadataProviderTest.java | 3 +-- 16 files changed, 16 insertions(+), 30 deletions(-) diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraGenericRecordSink.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraGenericRecordSink.java index 3ecc2ee4d575a..1798416517f99 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraGenericRecordSink.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraGenericRecordSink.java @@ -1,4 +1,4 @@ -/** +/* * 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 @@ -16,7 +16,6 @@ * specific language governing permissions and limitations * under the License. */ - package org.apache.pulsar.io.cassandra; import org.apache.pulsar.client.api.schema.GenericRecord; diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/BoundStatementProvider.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/BoundStatementProvider.java index 1b0b03ba0dbe1..f56e295548001 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/BoundStatementProvider.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/BoundStatementProvider.java @@ -1,4 +1,4 @@ -/** +/* * 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 @@ -16,7 +16,6 @@ * specific language governing permissions and limitations * under the License. */ - package org.apache.pulsar.io.cassandra.util; import com.datastax.driver.core.BoundStatement; diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/CassandraConnector.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/CassandraConnector.java index 20b8cfd200194..373404bca5e16 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/CassandraConnector.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/CassandraConnector.java @@ -1,4 +1,4 @@ -/** +/* * 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 @@ -16,7 +16,6 @@ * specific language governing permissions and limitations * under the License. */ - package org.apache.pulsar.io.cassandra.util; import com.datastax.driver.core.Cluster; diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/GenericRecordWrapper.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/GenericRecordWrapper.java index 3e1b24b16ff20..f7f725be366e6 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/GenericRecordWrapper.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/GenericRecordWrapper.java @@ -1,4 +1,4 @@ -/** +/* * 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 @@ -16,7 +16,6 @@ * specific language governing permissions and limitations * under the License. */ - package org.apache.pulsar.io.cassandra.util; import org.apache.pulsar.client.api.schema.GenericRecord; diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/RecordWrapper.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/RecordWrapper.java index b6a1b2b5bc429..628f0c430f1dc 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/RecordWrapper.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/RecordWrapper.java @@ -1,4 +1,4 @@ -/** +/* * 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 @@ -16,7 +16,6 @@ * specific language governing permissions and limitations * under the License. */ - package org.apache.pulsar.io.cassandra.util; import org.apache.commons.beanutils.converters.IntegerConverter; diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/StringRecordWrapper.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/StringRecordWrapper.java index 260d752ba2f1c..e7711ad72fa52 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/StringRecordWrapper.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/StringRecordWrapper.java @@ -1,4 +1,4 @@ -/** +/* * 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 @@ -16,7 +16,6 @@ * specific language governing permissions and limitations * under the License. */ - package org.apache.pulsar.io.cassandra.util; import com.google.gson.Gson; diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/TableMetadataProvider.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/TableMetadataProvider.java index 955202a060961..d978dcf31396a 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/TableMetadataProvider.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/TableMetadataProvider.java @@ -1,4 +1,4 @@ -/** +/* * 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 @@ -16,7 +16,6 @@ * specific language governing permissions and limitations * under the License. */ - package org.apache.pulsar.io.cassandra.util; import com.datastax.driver.core.ColumnMetadata; diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/package-info.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/package-info.java index 61699fd875df2..904579e812163 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/package-info.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/package-info.java @@ -1,4 +1,4 @@ -/** +/* * 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 diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java index ee129e74107b7..93628712d5537 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java @@ -1,4 +1,4 @@ -/** +/* * 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 @@ -16,7 +16,6 @@ * specific language governing permissions and limitations * under the License. */ - package org.apache.pulsar.io.cassandra; import java.io.File; diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/AbstractGenericRecordProducer.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/AbstractGenericRecordProducer.java index 2494756844b1f..1ce385c0e4c96 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/AbstractGenericRecordProducer.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/AbstractGenericRecordProducer.java @@ -1,4 +1,4 @@ -/** +/* * 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 @@ -16,7 +16,6 @@ * specific language governing permissions and limitations * under the License. */ - package org.apache.pulsar.io.cassandra.producers; import org.apache.pulsar.client.api.Schema; diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicProducerThread.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicProducerThread.java index a80f4d2cf326a..4c71ceeca9a4a 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicProducerThread.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicProducerThread.java @@ -1,4 +1,4 @@ -/** +/* * 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 @@ -16,7 +16,6 @@ * specific language governing permissions and limitations * under the License. */ - package org.apache.pulsar.io.cassandra.producers; import lombok.extern.slf4j.Slf4j; diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java index 2282de98249e7..f8c282574fba8 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java @@ -1,4 +1,4 @@ -/** +/* * 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 @@ -16,7 +16,6 @@ * specific language governing permissions and limitations * under the License. */ - package org.apache.pulsar.io.cassandra.producers; import com.google.gson.Gson; diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ObservationSchemaRecordProducer.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ObservationSchemaRecordProducer.java index 60cf753a42a1e..e77ebf72b70b3 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ObservationSchemaRecordProducer.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ObservationSchemaRecordProducer.java @@ -1,4 +1,4 @@ -/** +/* * 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 @@ -16,7 +16,6 @@ * specific language governing permissions and limitations * under the License. */ - package org.apache.pulsar.io.cassandra.producers; import org.apache.pulsar.client.api.Schema; diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ReadingSchemaRecordProducer.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ReadingSchemaRecordProducer.java index 126c8312de61b..e13309cdfa4aa 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ReadingSchemaRecordProducer.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ReadingSchemaRecordProducer.java @@ -1,4 +1,4 @@ -/** +/* * 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 @@ -16,7 +16,6 @@ * specific language governing permissions and limitations * under the License. */ - package org.apache.pulsar.io.cassandra.producers; import org.apache.pulsar.client.api.Schema; diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java index 6dde31dd67ad7..ab6f712206ab5 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java @@ -1,4 +1,4 @@ -/** +/* * 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 diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java index f14f78473eeff..1f37f740bebe0 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java @@ -1,4 +1,4 @@ -/** +/* * 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 @@ -16,7 +16,6 @@ * specific language governing permissions and limitations * under the License. */ - package org.apache.pulsar.io.cassandra.util; import org.apache.pulsar.io.cassandra.CassandraSinkConfig; From 21cedbabdb0a4207d1e08e4659530f0525c2f546 Mon Sep 17 00:00:00 2001 From: david-streamlio <35466513+david-streamlio@users.noreply.github.com> Date: Thu, 15 Dec 2022 16:26:48 -0800 Subject: [PATCH 06/14] Switched to TestNG framework --- .../producers/AbstractGenericRecordProducer.java | 2 +- .../producers/InputTopicStringProducer.java | 9 ++++----- .../io/cassandra/util/CassandraConnectorTest.java | 15 +++++---------- .../cassandra/util/TableMetadataProviderTest.java | 11 +++++------ 4 files changed, 15 insertions(+), 22 deletions(-) diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/AbstractGenericRecordProducer.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/AbstractGenericRecordProducer.java index 1ce385c0e4c96..96c341c46f125 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/AbstractGenericRecordProducer.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/AbstractGenericRecordProducer.java @@ -19,7 +19,7 @@ package org.apache.pulsar.io.cassandra.producers; import org.apache.pulsar.client.api.Schema; -import org.apache.pulsar.client.api.schema.*; +import org.apache.pulsar.client.api.schema.GenericRecord; import org.apache.pulsar.common.schema.SchemaInfo; @SuppressWarnings({"unchecked", "rawtypes"}) diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java index f8c282574fba8..b5e0d7c2ea010 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java @@ -31,7 +31,7 @@ @SuppressWarnings({"unchecked", "rawtypes"}) public class InputTopicStringProducer extends InputTopicProducerThread { - private Random rnd = new Random(); + private final Random rnd = new Random(); int lastReadingId = rnd.nextInt(900000); public InputTopicStringProducer(String brokerUrl, String inputTopic) { @@ -57,13 +57,12 @@ String getValue() { elements.put("reporting_area", lastReadingId + ""); elements.put("hour_observed", rnd.nextInt(24)); elements.put("date_observed", "2022-06-18"); - elements.put("latitude", Float.valueOf(40.021f)); - elements.put("longitude", Float.valueOf(-122.33f)); + elements.put("latitude", 40.021f); + elements.put("longitude", -122.33f); Gson gson = new Gson(); Type gsonType = new TypeToken(){}.getType(); - String gsonString = gson.toJson(elements,gsonType); - return gsonString; + return gson.toJson(elements,gsonType); } @Override diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java index ab6f712206ab5..78eabb8042b40 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java @@ -19,19 +19,16 @@ package org.apache.pulsar.io.cassandra.util; import org.apache.pulsar.io.cassandra.CassandraSinkConfig; -import org.apache.pulsar.io.cassandra.util.CassandraConnector; -import org.junit.Ignore; -import org.junit.Test; +import org.testng.annotations.Test; -import static org.junit.Assert.assertNotNull; import static org.testng.AssertJUnit.assertEquals; +import static org.testng.AssertJUnit.assertNotNull; public class CassandraConnectorTest { private CassandraSinkConfig config; - @Test - @Ignore + @Test(enabled = false) public final void securedTest() { config = new CassandraSinkConfig(); config.setRoots("localhost"); @@ -43,8 +40,7 @@ public final void securedTest() { assertNotNull(connector.getSession()); } - @Test - @Ignore + @Test(enabled = false) public final void getObservationPreparedStatementTest() { config = new CassandraSinkConfig(); config.setRoots("localhost"); @@ -57,8 +53,7 @@ public final void getObservationPreparedStatementTest() { assertEquals("INSERT INTO airquality.observation (key, observed) VALUES (?, ?)", connector.getPreparedStatement().getQueryString()); } - @Test - @Ignore + @Test(enabled = false) public final void getReadingPreparedStatementTest() { config = new CassandraSinkConfig(); config.setRoots("localhost"); diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java index 1f37f740bebe0..135634c1a74ba 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java @@ -19,16 +19,15 @@ package org.apache.pulsar.io.cassandra.util; import org.apache.pulsar.io.cassandra.CassandraSinkConfig; -import org.junit.Ignore; -import org.junit.Test; -import static org.junit.Assert.assertEquals; -import static org.junit.Assert.assertNotNull; +import org.testng.annotations.Test; + +import static org.testng.Assert.assertNotNull; +import static org.testng.AssertJUnit.assertEquals; public class TableMetadataProviderTest { - @Test - @Ignore + @Test(enabled = false) public final void getTableDefinitionTest() { CassandraSinkConfig config = new CassandraSinkConfig(); From fd370b4e697598bd8f6981af358dd9db13e8c078 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Fri, 6 Feb 2026 15:09:10 +0200 Subject: [PATCH 07/14] Revert unnecessary changes to pom.xml files that ignore cassandra.version property --- pom.xml | 6 +++++- tests/integration/pom.xml | 3 +-- 2 files changed, 6 insertions(+), 3 deletions(-) diff --git a/pom.xml b/pom.xml index bea090abb3899..3cc20802fa98f 100644 --- a/pom.xml +++ b/pom.xml @@ -1422,7 +1422,11 @@ flexible messaging model and an intuitive client API. pom import - + + com.datastax.cassandra + cassandra-driver-core + ${cassandra.version} + org.assertj assertj-core diff --git a/tests/integration/pom.xml b/tests/integration/pom.xml index 62d3854810064..ee940a36e7c1c 100644 --- a/tests/integration/pom.xml +++ b/tests/integration/pom.xml @@ -98,8 +98,7 @@ com.datastax.cassandra cassandra-driver-core - 3.11.2 - + io.netty netty-handler From 0f38c18bc6b6e628281ed32703c097658f340409 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Fri, 6 Feb 2026 15:13:06 +0200 Subject: [PATCH 08/14] Fix checkstyle --- .../pulsar/io/cassandra/CassandraSinkExec.java | 3 +-- .../producers/InputTopicProducerThread.java | 3 +-- .../producers/InputTopicStringProducer.java | 7 +++---- .../producers/ReadingSchemaRecordProducer.java | 3 +-- .../cassandra/util/CassandraConnectorTest.java | 16 ++++++++-------- .../util/TableMetadataProviderTest.java | 6 ++---- 6 files changed, 16 insertions(+), 22 deletions(-) diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java index 93628712d5537..613c21ad0cbf5 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/CassandraSinkExec.java @@ -25,7 +25,6 @@ import java.util.HashMap; import java.util.Map; import java.util.concurrent.TimeUnit; - import org.apache.pulsar.common.io.SinkConfig; import org.apache.pulsar.functions.LocalRunner; import org.apache.pulsar.io.cassandra.producers.InputTopicProducerThread; @@ -92,7 +91,7 @@ private static Map getConfigs() throws IOException { return configs; } - private static final void sendData() throws InterruptedException { + private static void sendData() throws InterruptedException { TimeUnit.SECONDS.sleep(10); InputTopicProducerThread writer = new ReadingSchemaRecordProducer(BROKER_URL, INPUT_TOPIC); writer.run(); diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicProducerThread.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicProducerThread.java index 4c71ceeca9a4a..260cbbc1db602 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicProducerThread.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicProducerThread.java @@ -18,14 +18,13 @@ */ package org.apache.pulsar.io.cassandra.producers; +import java.util.Random; import lombok.extern.slf4j.Slf4j; import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.PulsarClient; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.Schema; -import java.util.Random; - @SuppressWarnings({"unchecked", "rawtypes"}) @Slf4j public abstract class InputTopicProducerThread implements Runnable { diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java index b5e0d7c2ea010..7e2ef76d1fd08 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java @@ -20,20 +20,19 @@ import com.google.gson.Gson; import com.google.gson.reflect.TypeToken; -import org.apache.pulsar.client.api.Schema; - import java.lang.reflect.Type; import java.util.HashMap; import java.util.Random; import java.util.SortedMap; import java.util.TreeMap; +import org.apache.pulsar.client.api.Schema; @SuppressWarnings({"unchecked", "rawtypes"}) public class InputTopicStringProducer extends InputTopicProducerThread { private final Random rnd = new Random(); int lastReadingId = rnd.nextInt(900000); - + public InputTopicStringProducer(String brokerUrl, String inputTopic) { super(brokerUrl, inputTopic); } @@ -62,7 +61,7 @@ String getValue() { Gson gson = new Gson(); Type gsonType = new TypeToken(){}.getType(); - return gson.toJson(elements,gsonType); + return gson.toJson(elements, gsonType); } @Override diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ReadingSchemaRecordProducer.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ReadingSchemaRecordProducer.java index e13309cdfa4aa..e533a32ef6b23 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ReadingSchemaRecordProducer.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/ReadingSchemaRecordProducer.java @@ -18,6 +18,7 @@ */ package org.apache.pulsar.io.cassandra.producers; +import java.util.Random; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.schema.GenericRecord; import org.apache.pulsar.client.api.schema.RecordSchemaBuilder; @@ -25,8 +26,6 @@ import org.apache.pulsar.common.schema.SchemaInfo; import org.apache.pulsar.common.schema.SchemaType; -import java.util.Random; - @SuppressWarnings({"unchecked", "rawtypes"}) public class ReadingSchemaRecordProducer extends AbstractGenericRecordProducer { diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java index 78eabb8042b40..4a7d65cac965f 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java @@ -18,11 +18,10 @@ */ package org.apache.pulsar.io.cassandra.util; -import org.apache.pulsar.io.cassandra.CassandraSinkConfig; -import org.testng.annotations.Test; - import static org.testng.AssertJUnit.assertEquals; import static org.testng.AssertJUnit.assertNotNull; +import org.apache.pulsar.io.cassandra.CassandraSinkConfig; +import org.testng.annotations.Test; public class CassandraConnectorTest { @@ -50,7 +49,8 @@ public final void getObservationPreparedStatementTest() { config.setKeyspace("airquality"); CassandraConnector connector = new CassandraConnector(config); - assertEquals("INSERT INTO airquality.observation (key, observed) VALUES (?, ?)", connector.getPreparedStatement().getQueryString()); + assertEquals("INSERT INTO airquality.observation (key, observed) VALUES (?, ?)", + connector.getPreparedStatement().getQueryString()); } @Test(enabled = false) @@ -63,10 +63,10 @@ public final void getReadingPreparedStatementTest() { config.setKeyspace("airquality"); CassandraConnector connector = new CassandraConnector(config); - assertEquals("INSERT INTO airquality.reading " + - "(reporting_area, avg_ozone, avg_pm10, avg_pm25, date_observed, hour_observed, latitude, " + - "local_time_zone, longitude, max_ozone, max_pm10, max_pm25, min_ozone, min_pm10, min_pm25, " + - "readingid, state_code) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", + assertEquals("INSERT INTO airquality.reading " + + "(reporting_area, avg_ozone, avg_pm10, avg_pm25, date_observed, hour_observed, latitude, " + + "local_time_zone, longitude, max_ozone, max_pm10, max_pm25, min_ozone, min_pm10, min_pm25, " + + "readingid, state_code) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", connector.getPreparedStatement().getQueryString()); } } diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java index 135634c1a74ba..0921d4af6e035 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java @@ -18,12 +18,10 @@ */ package org.apache.pulsar.io.cassandra.util; -import org.apache.pulsar.io.cassandra.CassandraSinkConfig; - -import org.testng.annotations.Test; - import static org.testng.Assert.assertNotNull; import static org.testng.AssertJUnit.assertEquals; +import org.apache.pulsar.io.cassandra.CassandraSinkConfig; +import org.testng.annotations.Test; public class TableMetadataProviderTest { From 975cf7d6b77337b0417d81c2c54944382e2d2872 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Fri, 6 Feb 2026 15:36:41 +0200 Subject: [PATCH 09/14] Replace Gson usage with Jackson --- .../io/cassandra/util/StringRecordWrapper.java | 16 +++++++++------- .../producers/InputTopicStringProducer.java | 16 +++++++++------- 2 files changed, 18 insertions(+), 14 deletions(-) diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/StringRecordWrapper.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/StringRecordWrapper.java index e7711ad72fa52..6ca16e6d54069 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/StringRecordWrapper.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/StringRecordWrapper.java @@ -18,20 +18,22 @@ */ package org.apache.pulsar.io.cassandra.util; -import com.google.gson.Gson; -import com.google.gson.reflect.TypeToken; -import java.util.HashMap; +import com.fasterxml.jackson.databind.ObjectMapper; +import java.io.IOException; +import java.io.UncheckedIOException; import java.util.Map; public class StringRecordWrapper extends RecordWrapper { - - private static final Gson gson = new Gson(); + private static final ObjectMapper MAPPER = new ObjectMapper(); private Map valuesMap; public StringRecordWrapper(String jsonString) { super(jsonString); - valuesMap = gson.fromJson(jsonString, - new TypeToken>() {}.getType()); + try { + valuesMap = MAPPER.readValue(jsonString, Map.class); + } catch (IOException e) { + throw new UncheckedIOException(e); + } } @Override diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java index 7e2ef76d1fd08..7ada099728ea0 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/producers/InputTopicStringProducer.java @@ -18,10 +18,9 @@ */ package org.apache.pulsar.io.cassandra.producers; -import com.google.gson.Gson; -import com.google.gson.reflect.TypeToken; -import java.lang.reflect.Type; -import java.util.HashMap; +import com.fasterxml.jackson.core.JsonProcessingException; +import com.fasterxml.jackson.databind.ObjectMapper; +import java.io.UncheckedIOException; import java.util.Random; import java.util.SortedMap; import java.util.TreeMap; @@ -59,9 +58,12 @@ String getValue() { elements.put("latitude", 40.021f); elements.put("longitude", -122.33f); - Gson gson = new Gson(); - Type gsonType = new TypeToken(){}.getType(); - return gson.toJson(elements, gsonType); + ObjectMapper objectMapper = new ObjectMapper(); + try { + return objectMapper.writeValueAsString(elements); + } catch (JsonProcessingException e) { + throw new UncheckedIOException(e); + } } @Override From a7bf2e2371b0c23cc013778c5078516a12d0ec43 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Fri, 6 Feb 2026 15:40:09 +0200 Subject: [PATCH 10/14] Remove version for beanutils --- pulsar-io/cassandra/pom.xml | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-io/cassandra/pom.xml b/pulsar-io/cassandra/pom.xml index 3b365bf478f64..caf75f60747a1 100644 --- a/pulsar-io/cassandra/pom.xml +++ b/pulsar-io/cassandra/pom.xml @@ -77,7 +77,6 @@ commons-beanutils commons-beanutils - 1.9.4 compile From 071633c9c839dc234ff8dc95613312c03b820316 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Fri, 6 Feb 2026 16:14:13 +0200 Subject: [PATCH 11/14] Remove unnecessary dependency --- .../pulsar/io/cassandra/util/BoundStatementProvider.java | 5 ----- 1 file changed, 5 deletions(-) diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/BoundStatementProvider.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/BoundStatementProvider.java index f56e295548001..983bedf7564ed 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/BoundStatementProvider.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/BoundStatementProvider.java @@ -20,14 +20,9 @@ import com.datastax.driver.core.BoundStatement; import com.datastax.driver.core.PreparedStatement; -import org.apache.commons.beanutils.converters.IntegerConverter; -import org.apache.commons.beanutils.converters.NumberConverter; @SuppressWarnings("rawtypes") public class BoundStatementProvider { - - NumberConverter converter = new IntegerConverter(); - final TableMetadataProvider.TableDefinition tableDefinition; public BoundStatementProvider(TableMetadataProvider.TableDefinition tableDefinition) { From e09e03f770673ad862e34f7fef7f818d9e06e6a5 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Fri, 6 Feb 2026 16:14:29 +0200 Subject: [PATCH 12/14] Add Testcontainers dependency --- pulsar-io/cassandra/pom.xml | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/pulsar-io/cassandra/pom.xml b/pulsar-io/cassandra/pom.xml index caf75f60747a1..2ab35a0c03649 100644 --- a/pulsar-io/cassandra/pom.xml +++ b/pulsar-io/cassandra/pom.xml @@ -80,6 +80,12 @@ compile + + org.testcontainers + cassandra + test + + From 71191c7b7f0e3ba3af73e54d8209d8bcea119aef Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Fri, 6 Feb 2026 16:38:08 +0200 Subject: [PATCH 13/14] Use real Cassandra with Testcontainers --- .../io/cassandra/util/CassandraConnector.java | 2 +- .../cassandra/util/AbstractCassandraTest.java | 53 ++++++++++++++++ .../util/CassandraConnectorTest.java | 32 ++++------ .../util/TableMetadataProviderTest.java | 14 ++--- .../cassandra/src/test/resources/init.cql | 60 +++++++++++++++++++ 5 files changed, 132 insertions(+), 29 deletions(-) create mode 100644 pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/AbstractCassandraTest.java create mode 100644 pulsar-io/cassandra/src/test/resources/init.cql diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/CassandraConnector.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/CassandraConnector.java index 373404bca5e16..ad9eb329b2174 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/CassandraConnector.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/util/CassandraConnector.java @@ -28,7 +28,7 @@ import java.util.List; import org.apache.pulsar.io.cassandra.CassandraSinkConfig; -public class CassandraConnector { +public class CassandraConnector implements AutoCloseable { private Cluster cluster; private Session session; diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/AbstractCassandraTest.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/AbstractCassandraTest.java new file mode 100644 index 0000000000000..cf88a334e1d65 --- /dev/null +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/AbstractCassandraTest.java @@ -0,0 +1,53 @@ +/* + * 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.pulsar.io.cassandra.util; + +import java.net.InetSocketAddress; +import org.apache.pulsar.io.cassandra.CassandraSinkConfig; +import org.testcontainers.cassandra.CassandraContainer; +import org.testng.annotations.AfterClass; +import org.testng.annotations.BeforeClass; + +public class AbstractCassandraTest { + protected CassandraSinkConfig config; + protected CassandraContainer cassandraContainer; + + @BeforeClass + public void startCassandraContainer() { + cassandraContainer = new CassandraContainer("cassandra:3.11"); + cassandraContainer.withInitScript("init.cql"); + cassandraContainer.start(); + } + + @AfterClass(alwaysRun = true) + public void stopCassandraContainer() { + if (cassandraContainer != null) { + cassandraContainer.stop(); + cassandraContainer = null; + } + } + + protected void createSinkConfig() { + config = new CassandraSinkConfig(); + InetSocketAddress contactPoint = cassandraContainer.getContactPoint(); + config.setRoots(contactPoint.getHostString() + ":" + contactPoint.getPort()); + config.setUserName(cassandraContainer.getUsername()); + config.setPassword(cassandraContainer.getPassword()); + } +} diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java index 4a7d65cac965f..079f112180d67 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/CassandraConnectorTest.java @@ -20,51 +20,43 @@ import static org.testng.AssertJUnit.assertEquals; import static org.testng.AssertJUnit.assertNotNull; -import org.apache.pulsar.io.cassandra.CassandraSinkConfig; +import lombok.Cleanup; import org.testng.annotations.Test; -public class CassandraConnectorTest { +public class CassandraConnectorTest extends AbstractCassandraTest { - private CassandraSinkConfig config; - - @Test(enabled = false) + @Test public final void securedTest() { - config = new CassandraSinkConfig(); - config.setRoots("localhost"); - config.setUserName("cassandra"); - config.setPassword("cassandra"); + createSinkConfig(); + @Cleanup CassandraConnector connector = new CassandraConnector(config); connector.connect(); assertNotNull(connector.getSession()); } - @Test(enabled = false) + @Test public final void getObservationPreparedStatementTest() { - config = new CassandraSinkConfig(); - config.setRoots("localhost"); - config.setUserName("cassandra"); - config.setPassword("cassandra"); + createSinkConfig(); config.setColumnFamily("observation"); config.setKeyspace("airquality"); + @Cleanup CassandraConnector connector = new CassandraConnector(config); assertEquals("INSERT INTO airquality.observation (key, observed) VALUES (?, ?)", connector.getPreparedStatement().getQueryString()); } - @Test(enabled = false) + @Test public final void getReadingPreparedStatementTest() { - config = new CassandraSinkConfig(); - config.setRoots("localhost"); - config.setUserName("cassandra"); - config.setPassword("cassandra"); + createSinkConfig(); config.setColumnFamily("reading"); config.setKeyspace("airquality"); + @Cleanup CassandraConnector connector = new CassandraConnector(config); assertEquals("INSERT INTO airquality.reading " - + "(reporting_area, avg_ozone, avg_pm10, avg_pm25, date_observed, hour_observed, latitude, " + + "(reporting_area, date_observed, hour_observed, avg_ozone, avg_pm10, avg_pm25, latitude, " + "local_time_zone, longitude, max_ozone, max_pm10, max_pm25, min_ozone, min_pm10, min_pm25, " + "readingid, state_code) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)", connector.getPreparedStatement().getQueryString()); diff --git a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java index 0921d4af6e035..57ef5cb96024c 100644 --- a/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java +++ b/pulsar-io/cassandra/src/test/java/org/apache/pulsar/io/cassandra/util/TableMetadataProviderTest.java @@ -20,21 +20,19 @@ import static org.testng.Assert.assertNotNull; import static org.testng.AssertJUnit.assertEquals; -import org.apache.pulsar.io.cassandra.CassandraSinkConfig; +import lombok.Cleanup; import org.testng.annotations.Test; -public class TableMetadataProviderTest { +public class TableMetadataProviderTest extends AbstractCassandraTest { - @Test(enabled = false) + @Test public final void getTableDefinitionTest() { - CassandraSinkConfig config = new CassandraSinkConfig(); - config.setRoots("localhost"); - config.setUserName("cassandra"); - config.setPassword("cassandra"); + createSinkConfig(); config.setColumnFamily("observation"); config.setKeyspace("airquality"); + @Cleanup CassandraConnector connector = new CassandraConnector(config); TableMetadataProvider.TableDefinition table = @@ -45,6 +43,6 @@ public final void getTableDefinitionTest() { assertNotNull(table); assertEquals(17, table.getColumns().size()); assertEquals(1, table.getPartitionKeyColumns().size()); - assertEquals(1, table.getPrimaryKeyColumns().size()); + assertEquals(3, table.getPrimaryKeyColumns().size()); } } diff --git a/pulsar-io/cassandra/src/test/resources/init.cql b/pulsar-io/cassandra/src/test/resources/init.cql new file mode 100644 index 0000000000000..8ad23b3cd92a4 --- /dev/null +++ b/pulsar-io/cassandra/src/test/resources/init.cql @@ -0,0 +1,60 @@ +-- 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. + +-- Create keyspace for air quality data +CREATE KEYSPACE IF NOT EXISTS airquality +WITH replication = { + 'class': 'SimpleStrategy', + 'replication_factor': 1 +}; + +-- Use the airquality keyspace +USE airquality; + +-- Create observation table +CREATE TABLE IF NOT EXISTS observation ( + key text PRIMARY KEY, + observed text +); + +-- Create reading table with air quality measurements +CREATE TABLE IF NOT EXISTS reading ( + reporting_area text, + date_observed text, + hour_observed int, + readingid text, + avg_ozone double, + min_ozone double, + max_ozone double, + avg_pm10 double, + min_pm10 double, + max_pm10 double, + avg_pm25 double, + min_pm25 double, + max_pm25 double, + local_time_zone text, + state_code text, + latitude float, + longitude float, + PRIMARY KEY ((reporting_area), date_observed, hour_observed) +) WITH CLUSTERING ORDER BY (date_observed DESC, hour_observed DESC); + +-- Create index on readingid for queries by reading ID +CREATE INDEX IF NOT EXISTS reading_id_idx ON reading (readingid); + +-- Create index on state_code for queries by state +CREATE INDEX IF NOT EXISTS state_code_idx ON reading (state_code); From c3f18f9ab85abaefdf01bd55d61e807fc2d4fc59 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Fri, 6 Feb 2026 16:56:01 +0200 Subject: [PATCH 14/14] Revert renaming of CassandraAbstractSink --- .../{CassandraSink.java => CassandraAbstractSink.java} | 2 +- .../apache/pulsar/io/cassandra/CassandraGenericRecordSink.java | 2 +- .../org/apache/pulsar/io/cassandra/CassandraStringSink.java | 2 +- .../src/main/resources/META-INF/services/pulsar-io.yaml | 2 +- 4 files changed, 4 insertions(+), 4 deletions(-) rename pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/{CassandraSink.java => CassandraAbstractSink.java} (98%) diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSink.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraAbstractSink.java similarity index 98% rename from pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSink.java rename to pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraAbstractSink.java index d50430b4e2d2a..9615e4e9a569a 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraSink.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraAbstractSink.java @@ -42,7 +42,7 @@ type = IOType.SINK, help = "The CassandraStringSink is used for moving messages from Pulsar to Cassandra.", configClass = CassandraSinkConfig.class) -public abstract class CassandraSink implements Sink { +public abstract class CassandraAbstractSink implements Sink { CassandraConnector connector; CassandraSinkConfig cassandraSinkConfig; diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraGenericRecordSink.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraGenericRecordSink.java index 1798416517f99..d986efbc6fb2b 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraGenericRecordSink.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraGenericRecordSink.java @@ -23,7 +23,7 @@ import org.apache.pulsar.io.cassandra.util.GenericRecordWrapper; import org.apache.pulsar.io.cassandra.util.RecordWrapper; -public class CassandraGenericRecordSink extends CassandraSink { +public class CassandraGenericRecordSink extends CassandraAbstractSink { @Override RecordWrapper wrapRecord(Record record) { diff --git a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraStringSink.java b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraStringSink.java index 70c8f98f1ec39..ed90113d50f25 100644 --- a/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraStringSink.java +++ b/pulsar-io/cassandra/src/main/java/org/apache/pulsar/io/cassandra/CassandraStringSink.java @@ -22,7 +22,7 @@ import org.apache.pulsar.io.cassandra.util.RecordWrapper; import org.apache.pulsar.io.cassandra.util.StringRecordWrapper; -public class CassandraStringSink extends CassandraSink { +public class CassandraStringSink extends CassandraAbstractSink { @Override RecordWrapper wrapRecord(Record record) { diff --git a/pulsar-io/cassandra/src/main/resources/META-INF/services/pulsar-io.yaml b/pulsar-io/cassandra/src/main/resources/META-INF/services/pulsar-io.yaml index 8954d1621297a..b4863210c55a0 100644 --- a/pulsar-io/cassandra/src/main/resources/META-INF/services/pulsar-io.yaml +++ b/pulsar-io/cassandra/src/main/resources/META-INF/services/pulsar-io.yaml @@ -19,5 +19,5 @@ name: cassandra description: Writes data into Cassandra -sinkClass: org.apache.pulsar.io.cassandra.CassandraSink +sinkClass: org.apache.pulsar.io.cassandra.CassandraStringSink sinkConfigClass: org.apache.pulsar.io.cassandra.CassandraSinkConfig