Should we be able to configure threadlocal produce...
# general
j
Should we be able to configure threadlocal producers for Sources? The following is from my SourceConfig:
Copy code
.producerConfig(ProducerConfig.builder()
                        .maxPendingMessages(1000)
                        .compressionType(CompressionType.SNAPPY)
                        .useThreadLocalProducers(true)
                        .build())
When I run the Source I get a bunch of errors starting it:
Copy code
java.lang.RuntimeException: error writing discovered task to intermediate topic
	at org.apache.pulsar.functions.source.batch.BatchSourceExecutor.taskEater(BatchSourceExecutor.java:196) ~[org.apache.pulsar-pulsar-functions-instance-3.1.0.jar:3.1.0]
	at org.apache.pulsar.functions.source.batch.BatchSourceExecutor.lambda$triggerDiscover$1(BatchSourceExecutor.java:172) ~[org.apache.pulsar-pulsar-functions-instance-3.1.0.jar:3.1.0]
	at org.apache.pulsar.io.elasticsearch.ElasticSearchBatchSource.discover(ElasticSearchBatchSource.java:87) ~[?:?]
	at org.apache.pulsar.functions.source.batch.BatchSourceExecutor.lambda$triggerDiscover$2(BatchSourceExecutor.java:172) ~[org.apache.pulsar-pulsar-functions-instance-3.1.0.jar:3.1.0]
	at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) ~[?:?]
	at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) ~[?:?]
	at io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30) ~[io.netty-netty-common-4.1.94.Final.jar:4.1.94.Final]
	at java.lang.Thread.run(Thread.java:833) ~[?:?]"
Copy code
java.lang.RuntimeException: error writing discovered task to intermediate topic
	at org.apache.pulsar.functions.source.batch.BatchSourceExecutor.taskEater(BatchSourceExecutor.java:196) ~[org.apache.pulsar-pulsar-functions-instance-3.1.0.jar:3.1.0]
	at org.apache.pulsar.functions.source.batch.BatchSourceExecutor.lambda$triggerDiscover$1(BatchSourceExecutor.java:172) ~[org.apache.pulsar-pulsar-functions-instance-3.1.0.jar:3.1.0]
	at org.apache.pulsar.io.elasticsearch.ElasticSearchBatchSource.discover(ElasticSearchBatchSource.java:87) ~[?:?]
	at org.apache.pulsar.functions.source.batch.BatchSourceExecutor.lambda$triggerDiscover$2(BatchSourceExecutor.java:172) ~[org.apache.pulsar-pulsar-functions-instance-3.1.0.jar:3.1.0]
	at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1136) ~[?:?]
	at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:635) ~[?:?]
	at io.netty.util.concurrent.FastThreadLocalRunnable.run(FastThreadLocalRunnable.java:30) ~[io.netty-netty-common-4.1.94.Final.jar:4.1.94.Final]
	at java.lang.Thread.run(Thread.java:833) ~[?:?]
Copy code
java.lang.NullPointerException: Cannot invoke "java.util.Map.values()" because the return value of "java.lang.ThreadLocal.get()" is null
	at org.apache.pulsar.functions.instance.ContextImpl.close(ContextImpl.java:735) ~[org.apache.pulsar-pulsar-functions-instance-3.1.0.jar:3.1.0]
	at org.apache.pulsar.functions.instance.JavaInstance.close(JavaInstance.java:176) ~[org.apache.pulsar-pulsar-functions-instance-3.1.0.jar:3.1.0]
	at org.apache.pulsar.functions.instance.JavaInstanceRunnable.close(JavaInstanceRunnable.java:563) ~[org.apache.pulsar-pulsar-functions-instance-3.1.0.jar:3.1.0]
	at org.apache.pulsar.functions.instance.JavaInstanceRunnable.run(JavaInstanceRunnable.java:356) ~[org.apache.pulsar-pulsar-functions-instance-3.1.0.jar:3.1.0]
	at java.lang.Thread.run(Thread.java:833) ~[?:?]
The Source runs just fine when not setting the
useThreadLocal
parameter in the ProducerConfig.