This message was deleted.
# general
s
This message was deleted.
d
What is the result of the
schemas get
command AFTER you tried to manually create the topic?
o
I didn’t manually create the topic, I manually uploaded the schema
In that case I get:
Copy code
{
  "version": 2,
  "schemaInfo": {
    "name": "customer-processed",
    "schema": {
      "type": "record",
      "name": "User",
      "namespace": "data-src-conn",
      "fields": [
        {
          "name": "id",
          "type": [
            "null",
            "string"
          ]
        },
        {
          "name": "status",
          "type": [
            "null",
            "string"
          ]
        },
        {
          "name": "status_metadata",
          "type": [
            "string",
            "null"
          ]
        },
        {
          "name": "create_metadata",
          "type": [
            "string",
            "null"
          ]
        },
        {
          "name": "update_metadata",
          "type": [
            "string",
            "null"
          ]
        }
      ]
    },
    "type": "JSON",
    "timestamp": 1696426118370,
    "properties": {}
  }
}
I have attached the definition config file I used:
However, even after I do this, I still get the same error on the sink connector side
d
Thanks. Now that the schema exists on the topic, the next step is to configure the sink connector producer to use that schema. That can be done inside the yaml file for the connector
o
Can you let me know how exactly that can be done? From the configs in the connector, its not immediately clear to me how to do this
d
You can use the
customSchemaOutputs
property in the configs.
You will also need to set the
outputSchemaType
property to “JSON”. See the REST API for a detailed example.
o
Thank you @David K I will try this and let you know how it worked.
🤞 1
@David K could you clear a question up for me, do I pass these configs when creating the Function or when creating the Sink Connector?
If these are configs for the Pulsar Function then the limitation here is that they’re only supported for Java functions according to the docs, and not for Python functions
Possibly a rather silly question, Im trying to use a Java Pulsar Function now and since the message that Im trying to process is from the Debezium connector it has a Key Value schema. Now when I create the Java Pulsar Function, I get:
Copy code
Received error from server: Key schemas or Value schemas are different schema type, from key schema type is BYTES and to key schema is BYTES, from value schema is JSON and to value schema is BYTES
How can I resolve this?
d
What do the messages from the PostGRES Debezium Source Connector look like? Can you share the output? It looks like the message keys and values are both byte[]. Therefore your Pulsar function will need to consume messages using that schema. So in line 31 of your Python function code, you would first have to convert the input parameter to a String before passing it to the JSON parser.
If you are unsure of the data you are receiving, one approach that I like to use is the LocalRunner. This allows you to debug your Java function code against the actual topic data. This way, you can tweak your configuration to work as expected.
o
The schema for the Postgres source connector I have attached below, and a sample record using the Pulsar Client is:
Copy code
{
  "before": {
    "id": "f7a5d1df-7a7b-4ebe-8262-ab5078fb056f",
    "status": "INACTIVE",
    "status_metadata": "\"\"",
    "create_metadata": "{\"time\": 1696517856.219328, \"actor\": \"f7a5d1df-7a7b-4ebe-8262-ab5078fb056f\", \"actorType\": \"CUSTOMER\", \"callerService\": \"\"}",
    "update_metadata": "{\"time\": 1696517856.219328, \"actor\": \"f7a5d1df-7a7b-4ebe-8262-ab5078fb056f\", \"actorType\": \"CUSTOMER\", \"callerService\": \"\"}"
  },
  "after": null,
  "source": {
    "version": "1.9.7.Final",
    "connector": "postgresql",
    "name": "dbserver101",
    "ts_ms": 1696517858259,
    "snapshot": "false",
    "db": "customer",
    "sequence": "[\"36248032\",\"36248032\"]",
    "schema": "customer",
    "table": "customers",
    "txId": 1542,
    "lsn": 36248032,
    "xmin": null
  },
  "op": "d",
  "ts_ms": 1696517858479,
  "transaction": null
}
For all of the Pulsar resources, Im running them locally using the standalone distribution, which means for both the source connector and the Pulsar function I am also using
localrun
> So in line 31 of your Python function code, you would first have to convert the input parameter to a String before passing it to the JSON parser. Also Python Pulsar Functions dont support setting
outputSchemaType
and
customSchemaOutputs
. Thats why Im now writing a Java Pulsar Function, and Im using the ExclamationFunction example just to see how it works.
d
The output you shared showed a raw String, while the error message claimed the value was a byte array, so I am a bit confused. Can you share the output from the
pulsar-admin schema get
command for the topic?
o
Yes sorry for the confusion. There are two topics. the first is the input for the Pulsar Function and it has the schema that I shared above
The second is the output topic for the Pulsar Function and it has the following schema
Copy code
{
  "version": 1,
  "schemaInfo": {
    "name": "customer-processed",
    "schema": "\"string\"",
    "type": "AVRO",
    "timestamp": 1696519395927,
    "properties": {
      "__alwaysAllowNull": "true",
      "__jsr310ConversionEnabled": "false"
    }
  }
}
d
My set up is as follows: PostGRES Debezium Source Connector -> Input Topic -> Pulsar Function -> Output Topic -> DeltaLake Sink Connector
o
for both topics, i have used
pulsar-admin schema get
d
You cannot control the schema produced by the source connector. So you will have to configure your function to read the data as it is published. You CAN control how the data is written from your function, so you can define that as you see fit.
o
That makes sense, and that is my understanding as well. What Im missing is how to achieve that using the Pulsar Function.
d
No worries, we can work through it. 😃
🙌 1
I don’t see the output of the
schema get
command for the first topic. I only see the second one, which is an AVRO schema
o
d
Thanks. So the message keys are byte[], and the message payload/values are JSON schema. The next step is to configure your Java Function using
inputSpecs: - "input-topic-name": schemaType: JSON
, as defined here.
o
Copy code
bin/pulsar-admin functions localrun \
 --jar $PF_JAR_PATH \
 --className org.example.ExclamationFunction \
 --inputs <persistent://public/data-src-conn/dbserver101.customer.customers> \
 --output <persistent://public/data-src-conn/customer-processed> \
 --name ExclamationFunction \
 --schema-type json
This is already how Im running the function ^
I have restructured my command a little bit so that now Im running the following command
bin/pulsar-admin functions localrun --jar $PF_JAR_PATH --function-config-file $PF_CONFIG_PATH
and the yaml configs are as follows
Copy code
className: "org.example.ExclamationFunction"
inputs:
  - "<persistent://public/data-src-conn/dbserver101.customer.customers>"
output: "<persistent://public/data-src-conn/customer-processed>"
name: "ExclamationFunction"
schemaType: "json"
inputSpecs:
    schemaType: "json"
And Im still getting the same error as before
d
The inputSpecs property has to be defined as a map, with the topic name as the key, and then another map as the value.
Here is an example of how they are defined using a Localrunner. Note this works against a standalone pulsar instance running on localhost.
o
Alright I updated it to be this:
Copy code
className: "org.example.ExclamationFunction"
inputs:
  - "<persistent://public/data-src-conn/dbserver101.customer.customers>"
output: "<persistent://public/data-src-conn/customer-processed>"
name: "ExclamationFunction"
schemaType: "json"
inputSpecs:
    "<persistent://public/data-src-conn/dbserver101.customer.customers>":
        schemaType: "json"
which now leads me to the error:
Incompatible schema: exists schema type KEY_VALUE, new schema type JSON
🤔 1
d
Ah. That was my mistake. So the schema is
"type": "KEY_VALUE",
with the value being JSON.
So you will need to change your Function code to consume Key-Values.
Then in the LocalRunner code, do something like this;
Copy code
Schema<KeyValue<Integer, String>> kvSchema = Schema.KeyValue(
                Schema.BYTES,
                Schema.JSON(YOUR_CLASS),
                KeyValueEncodingType.SEPARATED
        );

        inputSpecs.put(IN, ConsumerConfig.builder().schemaType(
                kvSchema.getSchemaInfo().getType().toString()).build());
Then attach a debugger to see how that definition is created. This is a new one for me in terms of configuration in a YAML file. HTH.
o
Thanks @David K! When you say LocalRunner code, do you mean the Pulsar Function?
Is the LocalRunner achieving the same task as me using the admin cli
bin/pulsar-admin localrun
?
d
Yes, could you use the code I shared with you as an example? It connects to the Pulsar cluster to consumer the messages and then invokes your function code with each message. This allows you to attach a debugger, and step through the code, etc. You can see the message format, etc.
o
Ah I see, thanks once again David. I’ll report back with my findings 🙂
👍 1
@David K I was able to get a work around. I ended up using Unwrap Key-Value function provided by DataStax to get around the problem of consuming the KV Schema. Now my pipeline looks like this: Debezium PG Source Connector -> Topic 1 -> Unwrap KV Function -> Topic 2 -> Custom Transformation Function -> Topic 3 -> Deltalake Sink Connector For the custom function, I followed the JAVA Function examples, and created a custom object class
CustomerTableObject
and returned that to
Topic 3
. The schema for topic 3 now looks like:
Copy code
{
  "version": 0,
  "schemaInfo": {
    "name": "customer-processed-2",
    "schema": {
      "type": "record",
      "name": "CustomerTableObject",
      "namespace": "org.example",
      "fields": [
        {
          "name": "create_metadata",
          "type": [
            "null",
            "string"
          ]
        },
        {
          "name": "id",
          "type": [
            "null",
            "string"
          ]
        },
        {
          "name": "status",
          "type": [
            "null",
            "string"
          ]
        },
        {
          "name": "status_metadata",
          "type": [
            "null",
            "string"
          ]
        },
        {
          "name": "update_metadata",
          "type": [
            "null",
            "string"
          ]
        }
      ]
    },
    "type": "JSON",
    "timestamp": 1697034589493,
    "properties": {
      "__alwaysAllowNull": "true",
      "__jsr310ConversionEnabled": "false"
    }
  }
}
When I run the Sink Connector now, I can see the
_delta_log
directory being created and a parquet file of size 0 bytes which contains the schema. • When I query this delta table using Spark, I can only see the schema However I don’t see any of the payload being written down. When I send messages to topic 3 with the following payload:
Copy code
{
  "id": "apple",
  "status": "banana",
  "status_metadata": "caterpillar",
  "create_metadata": "dog",
  "update_metadata": "elep"
}
Nothing gets added. However, when I run
bin/pulsar-admin sinks status --tenant public --namespace data-src-conn --name delta_sink
I get this output:
Copy code
{
  "numInstances": 1,
  "numRunning": 1,
  "instances": [
    {
      "instanceId": 0,
      "status": {
        "running": true,
        "error": "",
        "numRestarts": 0,
        "numReadFromPulsar": 6,
        "numSystemExceptions": 0,
        "latestSystemExceptions": [],
        "numSinkExceptions": 0,
        "latestSinkExceptions": [],
        "numWrittenToSink": 6,
        "lastReceivedTime": 1697106408287,
        "workerId": "c-standalone-fw-localhost-8080"
      }
    }
  ]
}
I would expect
numReadFromPulsar
and/or
numWrittenToSink
to be 0 if its not successfully writing to my local file destination. Just for greater clarity, if I query the topic status using
bin/pulsar-admin topics stats <persistent://public/data-src-conn/topic-3>
I get the following:
Copy code
"public/data-src-conn/delta_sink" : {
      "msgRateOut" : 0.0,
      "msgThroughputOut" : 0.0,
      "bytesOutCounter" : 972,
      "msgOutCounter" : 6,
      "msgRateRedeliver" : 0.0,
      "messageAckRate" : 0.0,
      "chunkedMessageRate" : 0,
      "msgBacklog" : 3,
      "backlogSize" : 1035,
      "earliestMsgPublishTimeInBacklog" : 0,
      "msgBacklogNoDelayed" : 3,
      "blockedSubscriptionOnUnackedMsgs" : false,
      "msgDelayed" : 0,
      "unackedMessages" : 0,
      "type" : "Failover",
      "activeConsumerName" : "7d6e7",
      "msgRateExpired" : 0.0,
      "totalMsgExpired" : 0,
      "lastExpireTimestamp" : 0,
      "lastConsumedFlowTimestamp" : 1697106386342,
      "lastConsumedTimestamp" : 1697106408286,
      "lastAckedTimestamp" : 0,
      "lastMarkDeleteAdvancedTimestamp" : 0,
      "consumers" : [ {
        "msgRateOut" : 0.0,
        "msgThroughputOut" : 0.0,
        "bytesOutCounter" : 972,
        "msgOutCounter" : 6,
        "msgRateRedeliver" : 0.0,
        "messageAckRate" : 0.0,
        "chunkedMessageRate" : 0.0,
        "consumerName" : "7d6e7",
        "availablePermits" : 994,
        "unackedMessages" : 0,
        "avgMessagesPerEntry" : 2,
        "blockedConsumerOnUnackedMsgs" : false,
        "address" : "/127.0.0.1:57742",
        "connectedSince" : "2023-10-12T13:26:26.339217+03:00",
        "clientVersion" : "Pulsar-Java-v3.1.0",
        "lastAckedTimestamp" : 0,
        "lastConsumedTimestamp" : 1697106408286,
        "lastConsumedFlowTimestamp" : 1697106386342,
        "metadata" : {
          "instance_id" : "0",
          "application" : "pulsar-sink",
          "instance_hostname" : "Oneebs-MacBook-Pro.local",
          "id" : "public/data-src-conn/delta_sink"
        },
        "lastAckedTime" : "1970-01-01T04:00:00+04:00",
        "lastConsumedTime" : "2023-10-12T13:26:48.286+03:00"
      }
What am I missing? Why isn’t the Sink Connector writing delta parquet files to my local path?
@Hang Chen please have a look at this as well - facing the same issue as @Rajnesh
👍 1
d
Is the Sink connector possibly buffering messages internally before writing them to the local disk? Any errors inside the sink connector?
o
I cant see any error logs in the sink connector logs that would suggest something going wrong
@David K I have also tried following the example here https://stackoverflow.com/questions/73791829/delta-lake-sink-connector-for-apache-pulsar-with-minio-throws-java-lang-illegal and even then Im facing the same issue where the payload is not being written by the sink connector
👀 1
Any recommendations here @David K @Hang Chen?
h
@Yan Zhao