Currently if a user creates a sink and configures it to connect to a C* table that doesn't exist, the pod goes into a hung state where there are no additional logs, and pod metrics are not available.
17:59:24.377 [rockset/default/test-astra-db-0] INFO com.datastax.oss.driver.api.core.uuid.Uuids - PID obtained through native call to getpid(): 1
17:59:27.886 [rockset/default/test-astra-db-0] INFO com.datastax.oss.driver.internal.core.DefaultMavenCoordinates - DataStax Java driver for Apache Cassandra(R) (com.datastax.oss:java-driver-core) version 4.6.0
17:59:28.085 [rockset/default/test-astra-db-0] INFO com.datastax.oss.driver.internal.core.context.InternalDriverContext - Could not register Graph extensions; this is normal if Tinkerpop was explicitly excluded from classpath
17:59:28.288 [s1-admin-0] INFO com.datastax.oss.driver.internal.core.time.Clock - Using native clock for microsecond precision
17:59:29.375 [s1-io-0] INFO com.datastax.oss.driver.internal.core.channel.ChannelFactory - [s1] Failed to connect with protocol DSE_V2, retrying with DSE_V1
17:59:29.887 [s1-io-1] INFO com.datastax.oss.driver.internal.core.channel.ChannelFactory - [s1] Failed to connect with protocol DSE_V1, retrying with V4
17:59:32.990 [rockset/default/test-astra-db-0] ERROR com.datastax.oss.sink.pulsar.CassandraSinkTask - initialization error
com.datastax.oss.common.sink.ConfigException: Table test does not exist.
at com.datastax.oss.common.sink.state.LifeCycleManager.getTableMetadata(LifeCycleManager.java:375) ~[cassandra-sink-pulsar-1.4.0.jar:?]
at com.datastax.oss.common.sink.state.LifeCycleManager.lambda$buildInstanceState$8(LifeCycleManager.java:432) ~[cassandra-sink-pulsar-1.4.0.jar:?]
at java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:195) ~[?:?]
at java.util.ArrayList$ArrayListSpliterator.forEachRemaining(ArrayList.java:1655) ~[?:?]
at java.util.stream.AbstractPipeline.copyInto(AbstractPipeline.java:484) ~[?:?]
at java.util.stream.AbstractPipeline.wrapAndCopyInto(AbstractPipeline.java:474) ~[?:?]
at java.util.stream.ReduceOps$ReduceOp.evaluateSequential(ReduceOps.java:913) ~[?:?]
at java.util.stream.AbstractPipeline.evaluate(AbstractPipeline.java:234) ~[?:?]
at java.util.stream.ReferencePipeline.collect(ReferencePipeline.java:578) ~[?:?]
at com.datastax.oss.common.sink.state.LifeCycleManager.lambda$buildInstanceState$9(LifeCycleManager.java:453) ~[cassandra-sink-pulsar-1.4.0.jar:?]
at java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:195) ~[?:?]
at java.util.HashMap$ValueSpliterator.forEachRemaining(HashMap.java:1675) ~[?:?]
at java.util.stream.AbstractPipeline.copyInto(AbstractPipeline.java:484) ~[?:?]
at java.util.stream.AbstractPipeline.wrapAndCopyInto(AbstractPipeline.java:474) ~[?:?]
at java.util.stream.ReduceOps$ReduceOp.evaluateSequential(ReduceOps.java:913) ~[?:?]
at java.util.stream.AbstractPipeline.evaluate(AbstractPipeline.java:234) ~[?:?]
at java.util.stream.ReferencePipeline.collect(ReferencePipeline.java:578) ~[?:?]
at com.datastax.oss.common.sink.state.LifeCycleManager.buildInstanceState(LifeCycleManager.java:456) ~[cassandra-sink-pulsar-1.4.0.jar:?]
at com.datastax.oss.common.sink.state.LifeCycleManager.lambda$startTask$0(LifeCycleManager.java:112) ~[cassandra-sink-pulsar-1.4.0.jar:?]
at java.util.concurrent.ConcurrentHashMap.computeIfAbsent(ConcurrentHashMap.java:1705) ~[?:?]
at com.datastax.oss.common.sink.state.LifeCycleManager.startTask(LifeCycleManager.java:107) ~[cassandra-sink-pulsar-1.4.0.jar:?]
at com.datastax.oss.common.sink.AbstractSinkTask.start(AbstractSinkTask.java:67) ~[cassandra-sink-pulsar-1.4.0.jar:?]
at com.datastax.oss.sink.pulsar.CassandraSinkTask.open(CassandraSinkTask.java:143) ~[cassandra-sink-pulsar-1.4.0.jar:?]
at org.apache.pulsar.functions.instance.JavaInstanceRunnable.setupOutput(JavaInstanceRunnable.java:788) ~[?:?]
at org.apache.pulsar.functions.instance.JavaInstanceRunnable.setup(JavaInstanceRunnable.java:213) ~[?:?]
at org.apache.pulsar.functions.instance.JavaInstanceRunnable.run(JavaInstanceRunnable.java:235) ~[?:?]
at java.lang.Thread.run(Thread.java:829) ~[?:?]
17:59:32.993 [rockset/default/test-astra-db-0] INFO com.datastax.oss.common.sink.TaskStateManager - Task is stopped.
17:59:32.993 [rockset/default/test-astra-db-0] ERROR org.apache.pulsar.functions.instance.JavaInstanceRunnable - Sink open produced uncaught exception:
com.datastax.oss.common.sink.ConfigException: Table test does not exist.
at com.datastax.oss.common.sink.state.LifeCycleManager.getTableMetadata(LifeCycleManager.java:375) ~[cassandra-sink-pulsar-1.4.0.jar:?]
at com.datastax.oss.common.sink.state.LifeCycleManager.lambda$buildInstanceState$8(LifeCycleManager.java:432) ~[cassandra-sink-pulsar-1.4.0.jar:?]
at java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:195) ~[?:?]
at java.util.ArrayList$ArrayListSpliterator.forEachRemaining(ArrayList.java:1655) ~[?:?]
at java.util.stream.AbstractPipeline.copyInto(AbstractPipeline.java:484) ~[?:?]
at java.util.stream.AbstractPipeline.wrapAndCopyInto(AbstractPipeline.java:474) ~[?:?]
at java.util.stream.ReduceOps$ReduceOp.evaluateSequential(ReduceOps.java:913) ~[?:?]
at java.util.stream.AbstractPipeline.evaluate(AbstractPipeline.java:234) ~[?:?]
at java.util.stream.ReferencePipeline.collect(ReferencePipeline.java:578) ~[?:?]
at com.datastax.oss.common.sink.state.LifeCycleManager.lambda$buildInstanceState$9(LifeCycleManager.java:453) ~[cassandra-sink-pulsar-1.4.0.jar:?]
at java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:195) ~[?:?]
at java.util.HashMap$ValueSpliterator.forEachRemaining(HashMap.java:1675) ~[?:?]
at java.util.stream.AbstractPipeline.copyInto(AbstractPipeline.java:484) ~[?:?]
at java.util.stream.AbstractPipeline.wrapAndCopyInto(AbstractPipeline.java:474) ~[?:?]
at java.util.stream.ReduceOps$ReduceOp.evaluateSequential(ReduceOps.java:913) ~[?:?]
at java.util.stream.AbstractPipeline.evaluate(AbstractPipeline.java:234) ~[?:?]
at java.util.stream.ReferencePipeline.collect(ReferencePipeline.java:578) ~[?:?]
at com.datastax.oss.common.sink.state.LifeCycleManager.buildInstanceState(LifeCycleManager.java:456) ~[cassandra-sink-pulsar-1.4.0.jar:?]
at com.datastax.oss.common.sink.state.LifeCycleManager.lambda$startTask$0(LifeCycleManager.java:112) ~[cassandra-sink-pulsar-1.4.0.jar:?]
at java.util.concurrent.ConcurrentHashMap.computeIfAbsent(ConcurrentHashMap.java:1705) ~[?:?]
at com.datastax.oss.common.sink.state.LifeCycleManager.startTask(LifeCycleManager.java:107) ~[cassandra-sink-pulsar-1.4.0.jar:?]
at com.datastax.oss.common.sink.AbstractSinkTask.start(AbstractSinkTask.java:67) ~[cassandra-sink-pulsar-1.4.0.jar:?]
at com.datastax.oss.sink.pulsar.CassandraSinkTask.open(CassandraSinkTask.java:143) ~[cassandra-sink-pulsar-1.4.0.jar:?]
at org.apache.pulsar.functions.instance.JavaInstanceRunnable.setupOutput(JavaInstanceRunnable.java:788) ~[?:?]
at org.apache.pulsar.functions.instance.JavaInstanceRunnable.setup(JavaInstanceRunnable.java:213) ~[?:?]
at org.apache.pulsar.functions.instance.JavaInstanceRunnable.run(JavaInstanceRunnable.java:235) ~[?:?]
at java.lang.Thread.run(Thread.java:829) ~[?:?]
17:59:32.994 [rockset/default/test-astra-db-0] ERROR org.apache.pulsar.functions.instance.JavaInstanceRunnable - [rockset/default/test-astra-db:0] Uncaught exception in Java Instance
com.datastax.oss.common.sink.ConfigException: Table test does not exist.
at com.datastax.oss.common.sink.state.LifeCycleManager.getTableMetadata(LifeCycleManager.java:375) ~[?:?]
at com.datastax.oss.common.sink.state.LifeCycleManager.lambda$buildInstanceState$8(LifeCycleManager.java:432) ~[?:?]
at java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:195) ~[?:?]
at java.util.ArrayList$ArrayListSpliterator.forEachRemaining(ArrayList.java:1655) ~[?:?]
at java.util.stream.AbstractPipeline.copyInto(AbstractPipeline.java:484) ~[?:?]
at java.util.stream.AbstractPipeline.wrapAndCopyInto(AbstractPipeline.java:474) ~[?:?]
at java.util.stream.ReduceOps$ReduceOp.evaluateSequential(ReduceOps.java:913) ~[?:?]
at java.util.stream.AbstractPipeline.evaluate(AbstractPipeline.java:234) ~[?:?]
at java.util.stream.ReferencePipeline.collect(ReferencePipeline.java:578) ~[?:?]
at com.datastax.oss.common.sink.state.LifeCycleManager.lambda$buildInstanceState$9(LifeCycleManager.java:453) ~[?:?]
at java.util.stream.ReferencePipeline$3$1.accept(ReferencePipeline.java:195) ~[?:?]
at java.util.HashMap$ValueSpliterator.forEachRemaining(HashMap.java:1675) ~[?:?]
at java.util.stream.AbstractPipeline.copyInto(AbstractPipeline.java:484) ~[?:?]
at java.util.stream.AbstractPipeline.wrapAndCopyInto(AbstractPipeline.java:474) ~[?:?]
at java.util.stream.ReduceOps$ReduceOp.evaluateSequential(ReduceOps.java:913) ~[?:?]
at java.util.stream.AbstractPipeline.evaluate(AbstractPipeline.java:234) ~[?:?]
at java.util.stream.ReferencePipeline.collect(ReferencePipeline.java:578) ~[?:?]
at com.datastax.oss.common.sink.state.LifeCycleManager.buildInstanceState(LifeCycleManager.java:456) ~[?:?]
at com.datastax.oss.common.sink.state.LifeCycleManager.lambda$startTask$0(LifeCycleManager.java:112) ~[?:?]
at java.util.concurrent.ConcurrentHashMap.computeIfAbsent(ConcurrentHashMap.java:1705) ~[?:?]
at com.datastax.oss.common.sink.state.LifeCycleManager.startTask(LifeCycleManager.java:107) ~[?:?]
at com.datastax.oss.common.sink.AbstractSinkTask.start(AbstractSinkTask.java:67) ~[?:?]
at com.datastax.oss.sink.pulsar.CassandraSinkTask.open(CassandraSinkTask.java:143) ~[?:?]
at org.apache.pulsar.functions.instance.JavaInstanceRunnable.setupOutput(JavaInstanceRunnable.java:788) ~[org.apache.pulsar-pulsar-functions-instance-2.7.2.1.0.0.jar:2.7.2.1.0.0]
at org.apache.pulsar.functions.instance.JavaInstanceRunnable.setup(JavaInstanceRunnable.java:213) ~[org.apache.pulsar-pulsar-functions-instance-2.7.2.1.0.0.jar:2.7.2.1.0.0]
at org.apache.pulsar.functions.instance.JavaInstanceRunnable.run(JavaInstanceRunnable.java:235) ~[org.apache.pulsar-pulsar-functions-instance-2.7.2.1.0.0.jar:2.7.2.1.0.0]
at java.lang.Thread.run(Thread.java:829) ~[?:?]
17:59:33.073 [rockset/default/test-astra-db-0] INFO org.apache.pulsar.functions.instance.JavaInstanceRunnable - Closing instance
17:59:33.073 [rockset/default/test-astra-db-0] INFO com.datastax.oss.common.sink.TaskStateManager - Task is stopped.
17:59:33.076 [rockset/default/test-astra-db-0] INFO org.apache.pulsar.functions.instance.JavaInstanceRunnable - Unloading JAR files for function InstanceConfig(instanceId=0, functionId=722e3813-8c67-43f4-a91c-a7c7a0f7c948, functionVersion=ada96250-45f6-4292-9f23-315170fe12fd, functionDetails=tenant: "rockset"
namespace: "default"
name: "test-astra-db"
className: "org.apache.pulsar.functions.api.utils.IdentityFunction"
parallelism: 1
source {
typeClassName: "org.apache.pulsar.client.api.schema.GenericRecord"
timeoutMs: 5000
inputSpecs {
key: "persistent://rockset/default/test-topic"
value {
}
}
}
sink {
className: "com.datastax.oss.sink.pulsar.RecordCassandraSinkTask"
configs: "{\"auth\":{\"gssapi\":{\"service\":\"dse\"},\"password\":\"",\"provider\":\"None\",\"username\":\""},\"cloud.secureConnectBundle\":\"base64:",\"compression\":\"None\",\"connectionPoolLocalSize\":4,\"ignoreErrors\":\"None\",\"jmx\":true,\"maxConcurrentRequests\":500,\"maxNumberOfRecordsInBatch\":32,\"queryExecutionTimeout\":30,\"task.max\":1,\"tasks.max\":1,\"topic\":{\"test-topic\":{\"codec\":{\"date\":\"ISO_LOCAL_DATE\",\"locale\":\"en_US\",\"time\":\"ISO_LOCAL_TIME\",\"timeZone\":\"UTC\",\"timestamp\":\"CQL_TIMESTAMP\",\"unit\":\"MILLISECONDS\"},\"test\":{\"test\":{\"consistencyLevel\":\"LOCAL_ONE\",\"deletesEnabled\":true,\"mapping\":\"part\\u003dvalue.name, id\\u003dvalue.id, num\\u003dvalue.number, fact\\u003dvalue.isfact, added\\u003dnow()\",\"nullToUnset\":true,\"timestampTimeUnit\":\"MICROSECONDS\",\"ttl\":-1,\"ttlTimeUnit\":\"SECONDS\"}}}},\"topics\":\"test-topic\"}"
typeClassName: "org.apache.pulsar.client.api.schema.GenericRecord"
builtin: "cassandra-enhanced"
}
resources {
cpu: 0.5
ram: 524288000
disk: 1073741824
}
componentType: SINK
customRuntimeOptions: "{ \"nodeSelectorLabels\": { \"astra-node\": \"functionworker\"}, \"tolerations\" : [{ \"key\&q...
Currently if a user creates a sink and configures it to connect to a C* table that doesn't exist, the pod goes into a hung state where there are no additional logs, and pod metrics are not available.