This message was deleted.
# general
s
This message was deleted.
a
Can you confirm if the connector archive is uploaded correctly to the broker/worker?
o
Yes I can confirm. I am able to successfully run the connector using the
pulsar-admin
CLI commands
a
What does your setup look like, is the connector running in the broker or as a separate k8s pod? And are you running on plain Pulsar or function-mesh?
o
The connector runs as a separate k8s pod once its created. Im running this on plain pulsar - ver 2.11.2
a
Thanks for clarifying. Can you please try to add the
"className"
property to the root of the json?
I belive the value should be
<http://org.apache.pulsar.io|org.apache.pulsar.io>.kafka.connect.KafkaConnectSource
o
Im running into the same issue
FYI: I was able to run the connector using the same configs as before via
pulsar-admin source create..
command without providing
className
in the root of the json, and instead providing
kafkaConnectorSourceClass
in the
configs
field
a
Can you provide the working pulsar-admin command to cross-reference? That might help with finding the root cause
o
Sure, the command is:
./bin/pulsar-admin sources create --source-config-file /tmp/yugabyte-connector-config.yaml
when I create it via the command the connector config file exists in the broker pod within the
tmp
directory
a
Do you see any error in the logs of the connector pod when submitting the request?
o
no it works as expected, I can share those here?
a
I mean when you submit via REST API there could be a log indicating why it could not create the classloader for the connector
You can share the logs, just be aware of redacting e.g. credentials which could be in plaintext there
✅ 1
o
Copy code
2023-08-10T08:02:40,972    ERROR    SourcesImpl    Invalid register Source request @ /public/default/debezium-yb-connector-turbine-turbine-txnreqsessions
java.lang.IllegalArgumentException: Source package is not provided
	at org.apache.pulsar.functions.worker.rest.api.SourcesImpl.validateUpdateRequestParams(SourcesImpl.java:750)
	at org.apache.pulsar.functions.worker.rest.api.SourcesImpl.registerSource(SourcesImpl.java:167)
	at org.apache.pulsar.broker.admin.impl.SourcesBase.registerSource(SourcesBase.java:148)
	at jdk.internal.reflect.GeneratedMethodAccessor310.invoke(Unknown Source)
	at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
	at java.base/java.lang.reflect.Method.invoke(Method.java:568)
	at org.glassfish.jersey.server.model.internal.ResourceMethodInvocationHandlerFactory.lambda$static$0(ResourceMethodInvocationHandlerFactory.java:52)
	at org.glassfish.jersey.server.model.internal.AbstractJavaResourceMethodDispatcher$1.run(AbstractJavaResourceMethodDispatcher.java:124)
	at org.glassfish.jersey.server.model.internal.AbstractJavaResourceMethodDispatcher.invoke(AbstractJavaResourceMethodDispatcher.java:167)
	at org.glassfish.jersey.server.model.internal.JavaResourceMethodDispatcherProvider$VoidOutInvoker.doDispatch(JavaResourceMethodDispatcherProvider.java:159)
	at org.glassfish.jersey.server.model.internal.AbstractJavaResourceMethodDispatcher.dispatch(AbstractJavaResourceMethodDispatcher.java:79)
	at org.glassfish.jersey.server.model.ResourceMethodInvoker.invoke(ResourceMethodInvoker.java:475)
	at org.glassfish.jersey.server.model.ResourceMethodInvoker.apply(ResourceMethodInvoker.java:397)
	at org.glassfish.jersey.server.model.ResourceMethodInvoker.apply(ResourceMethodInvoker.java:81)
	at org.glassfish.jersey.server.ServerRuntime$1.run(ServerRuntime.java:255)
	at org.glassfish.jersey.internal.Errors$1.call(Errors.java:248)
	at org.glassfish.jersey.internal.Errors$1.call(Errors.java:244)
	at org.glassfish.jersey.internal.Errors.process(Errors.java:292)
	at org.glassfish.jersey.internal.Errors.process(Errors.java:274)
	at org.glassfish.jersey.internal.Errors.process(Errors.java:244)
	at org.glassfish.jersey.process.internal.RequestScope.runInScope(RequestScope.java:265)
	at org.glassfish.jersey.server.ServerRuntime.process(ServerRuntime.java:234)
	at org.glassfish.jersey.server.ApplicationHandler.handle(ApplicationHandler.java:680)
	at org.glassfish.jersey.servlet.WebComponent.serviceImpl(WebComponent.java:394)
	at org.glassfish.jersey.servlet.WebComponent.service(WebComponent.java:346)
@Alexander Preuß are these logs any help?
Can I get any help here?
a
@Oneeb I think something went wrong pasting the logs, I can’t see the Exception
o
@Alexander Preuß apologies, corrected it
a
It is an issue with the classloading in the worker. Only advice I can give is to double check once again if you can see the file on the worker node, then check if it has been correctly extracted in the workers NAR extraction directory and that the className is correct too
a
@Ming might have some insights into use of the REST API to create connectors
🙌 1
o
@Ming I would really appreciate help here if possible
m
@Oneeb this is an example of curl command I used on our cluster. You can ask GPT to convert Python code.
Copy code
$ curl -v --location --request POST '<https://pulsar-aws-useast2.api.streaming.datastax.com/admin/v3/sinks/mingt0/default/secondes-sink>' --header "Authorization: Bearer $TOKEN" --form 'sinkConfig="{\"archive\":\"builtin:\/\/elastic_search\",\"tenant\":\"mingt0\",\"namespace\":\"default\",\"name\":\"secondes-sink\",\"parallelism\":1,\"inputs\":[\"mingt0\/default\/astra\"],\"configs\":{\"elasticSearchUrl\":\"https:\/\/mingdev.es.us-east1.gcp.elastic-cloud.com:9243\",\"password\":\"removed...\",\"indexName\":\"devtest\"}}"'
o
@Ming where you’re pasing
sinkConfig
Im not passing anything. What should I pass here considering Im using a KCA Source?
a
@Oneeb should be your
sourceConfig
in place of his
sinkConfig
m
Yes. I used a sink example.
o
@Alexander Preuß @Ming coming back to try this now, if you see my first message here, Im already providing
sourceConfig
I have tried several variations of this and Im still getting the same error I’ve been getting as before
Source package is not provide
. Is there something Im missing here to create a Kafka Connect Adaptor source?
Stack trace:
Copy code
2023-08-30T21:21:41,945    ERROR    SourcesImpl    Invalid register Source request @ /public/default/debezium-yb-connector-turbine-turbine-txnreqsessions
java.lang.IllegalArgumentException: Source package is not provided
	at org.apache.pulsar.functions.worker.rest.api.SourcesImpl.validateUpdateRequestParams(SourcesImpl.java:692)
	at org.apache.pulsar.functions.worker.rest.api.SourcesImpl.registerSource(SourcesImpl.java:151)
	at org.apache.pulsar.broker.admin.impl.SourcesBase.registerSource(SourcesBase.java:144)
	at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
	at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:77)
	at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
	at java.base/java.lang.reflect.Method.invoke(Method.java:568)
	at org.glassfish.jersey.server.model.internal.ResourceMethodInvocationHandlerFactory.lambda$static$0(ResourceMethodInvocationHandlerFactory.java:52)
	at org.glassfish.jersey.server.model.internal.AbstractJavaResourceMethodDispatcher$1.run(AbstractJavaResourceMethodDispatcher.java:124)
	at org.glassfish.jersey.server.model.internal.AbstractJavaResourceMethodDispatcher.invoke(AbstractJavaResourceMethodDispatcher.java:167)
	at org.glassfish.jersey.server.model.internal.JavaResourceMethodDispatcherProvider$VoidOutInvoker.doDispatch(JavaResourceMethodDispatcherProvider.java:159)
	at org.glassfish.jersey.server.model.internal.AbstractJavaResourceMethodDispatcher.dispatch(AbstractJavaResourceMethodDispatcher.java:79)
	at org.glassfish.jersey.server.model.ResourceMethodInvoker.invoke(ResourceMethodInvoker.java:475)
	at org.glassfish.jersey.server.model.ResourceMethodInvoker.apply(ResourceMethodInvoker.java:397)
	at org.glassfish.jersey.server.model.ResourceMethodInvoker.apply(ResourceMethodInvoker.java:81)
	at org.glassfish.jersey.server.ServerRuntime$1.run(ServerRuntime.java:255)
	at org.glassfish.jersey.internal.Errors$1.call(Errors.java:248)
For Pulsar Admin Sources REST API Im still struggling with the
Source package is not provided
error. This is what my python script looks like:
Copy code
connector_config = {
            "tenant": self.tenant,
            "namespace": self.namespace,
            "topicName": self.topic_name,
            "archive": "connectors/pulsar-io-debezium-postgres-3.1.0.nar",
            "parallelism": 1,
            "processingGuarantees": "EXACTLY_ONCE",
            "configs": {
                "key.converter": "org.apache.kafka.connect.json.JsonConverter",
                "value.converter": "org.apache.kafka.connect.json.JsonConverter",
                "key.converter.schemas.enable": True,
                "value.converter.schemas.enable": True,
                "json-with-envelope": True,
                "database.hostname": self.db_hostname,
                "database.port": self.db_port,
                "database.user": self.db_access_key,
                "database.password": self.db_secret_key,
                "database.dbname": self.database_name,
                "database.server.name": f"{self.db_server_name}",
                "table.include.list": f"{self.schema_name}.{self.table_name}",
                "snapshot.mode": self.snapshot_mode,
                "pulsar.service.url": self.pulsar_broker_url,
                "plugin.name": "pgoutput"
            }
        }
        connector_config_json = json.dumps(connector_config)
        # print(connector_config)
        print(pulsar_admin_server)
        pulsar_connector_url = \
            f"{pulsar_admin_server}/admin/v3/sources/{self.tenant}/{self.namespace}/{self.connector_name}"

        mp_encoder = MultipartEncoder([("sourceConfig", (None, connector_config_json, "application/json"))])
        # mp_encoder = MultipartEncoder(fields={"sourceConfig": (None, connector_config_json, "application/json")})

        x = <http://requests.post|requests.post>(url=pulsar_connector_url,
                          data=mp_encoder, headers={"Content-Type": mp_encoder.content_type})
And the error that Im getting on the pulsar broker is:
Copy code
java.lang.IllegalArgumentException: Source package is not provided
	at org.apache.pulsar.functions.worker.rest.api.SourcesImpl.validateUpdateRequestParams(SourcesImpl.java:692) ~[org.apache.pulsar-pulsar-functions-worker-3.1.0.jar:3.1.0]
	at org.apache.pulsar.functions.worker.rest.api.SourcesImpl.registerSource(SourcesImpl.java:151) ~[org.apache.pulsar-pulsar-functions-worker-3.1.0.jar:3.1.0]
	at org.apache.pulsar.broker.admin.impl.SourcesBase.registerSource(SourcesBase.java:144) ~[org.apache.pulsar-pulsar-broker-3.1.0.jar:3.1.0]
	at jdk.internal.reflect.GeneratedMethodAccessor326.invoke(Unknown Source) ~[?:?]
	at jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) ~[?:?]
	at java.lang.reflect.Method.invoke(Method.java:568) ~[?:?]
	at org.glassfish.jersey.server.model.internal.ResourceMethodInvocationHandlerFactory.lambda$static$0(ResourceMethodInvocationHandlerFactory.java:52) ~[org.glassfish.jersey.core-jersey-server-2.34.jar:?]
	at org.glassfish.jersey.server.model.internal.AbstractJavaResourceMethodDispatcher$1.run(AbstractJavaResourceMethodDispatcher.java:124) ~[org.glassfish.jersey.core-jersey-server-2.34.jar:?]
	at org.glassfish.jersey.server.model.internal.AbstractJavaResourceMethodDispatcher.invoke(AbstractJavaResourceMethodDispatcher.java:167) ~[org.glassfish.jersey.core-jersey-server-2.34.jar:?]
	at org.glassfish.jersey.server.model.internal.JavaResourceMethodDispatcherProvider$VoidOutInvoker.doDispatch(JavaResourceMethodDispatcherProvider.java:159) ~[org.glassfish.jersey.core-jersey-server-2.34.jar:?]
	at org.glassfish.jersey.server.model.internal.AbstractJavaResourceMethodDispatcher.dispatch(AbstractJavaResourceMethodDispatcher.java:79) ~[org.glassfish.jersey.core-jersey-server-2.34.jar:?]
	at org.glassfish.jersey.server.model.ResourceMethodInvoker.invoke(ResourceMethodInvoker.java:475) ~[org.glassfish.jersey.core-jersey-server-2.34.jar:?]
Would really like to solve this issue now.
@Ming can I get your input here?
@Andrey can I be connected with someone who has worked on this?
m
Copy code
connectors/pulsar-io-debezium-postgres-3.1.0.nar
Have you tried
"archive": "<builtin://debezium-postgres>"
instead of the above?
I suppose you are using the pulsar-all which should have already included the debezium nar
o
Thank you very much Ming! This was the solution.
Can you point me to documentation which has this mentioned for builtin connectors and what to pass for each one of them?
m
Good to know that works for you. It should be in the manifest file in the nar or somewhere in the code.