Skip to content

C* sink should terminate if the configured table doesn't exist #14

Description

@pgier

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...

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions