Slackbot
09/21/2023, 10:59 PMNeng
09/22/2023, 12:44 AMOneeb
09/22/2023, 4:43 AMjson-with-envelope: False? Thats how I originally ran it and I would get the same error.Neng
09/22/2023, 4:54 AMOneeb
09/22/2023, 5:31 AMReceived error from server: org.apache.pulsar.broker.service.schema.exceptions.IncompatibleSchemaException: Key schemas or Value schemas are different schema type, from key schema type is BYTES and to key schema is JSON, from value schema is BYTES and to value schema is JSON caused by org.apache.pulsar.broker.service.schema.exceptions.IncompatibleSchemaException: Key schemas or Value schemas are different schema type, from key schema type is BYTES and to key schema is JSON, from value schema is BYTES and to value schema is JSONNeng
09/22/2023, 5:47 AMOneeb
09/22/2023, 6:19 AMOneeb
09/22/2023, 6:38 AMorg.apache.kafka.json.JsonConverter isnt having any effectOneeb
09/22/2023, 8:43 AMWe do not support KEY_VALUE. Any help here?Neng
09/22/2023, 7:44 PM{
"key": "",
"value": {
"type": "record",
"name": "Envelope",
"namespace": "mydbinstance.public.users",
"fields": [
{
"name": "before",
"type": [
"null",
{
"type": "record",
"name": "Value",
"fields": [
{
"name": "firstname",
"type": [
"null",
"string"
],
"default": null
},
{
"name": "lastname",
"type": [
"null",
"string"
],
"default": null
},
{
"name": "age",
"type": [
"null",
"int"
],
"default": null
}
],
"connect.name": "mydbinstance.public.users.Value"
}
],
"default": null
},
{
"name": "after",
"type": [
"null",
"Value"
],
"default": null
},
{
"name": "source",
"type": {
"type": "record",
"name": "Source",
"namespace": "io.debezium.connector.postgresql",
"fields": [
{
"name": "version",
"type": "string"
},
{
"name": "connector",
"type": "string"
},
{
"name": "name",
"type": "string"
},
{
"name": "ts_ms",
"type": "long"
},
{
"name": "snapshot",
"type": [
{
"type": "string",
"connect.version": 1,
"connect.parameters": {
"allowed": "true,last,false"
},
"connect.default": "false",
"connect.name": "io.debezium.data.Enum"
},
"null"
],
"default": "false"
},
{
"name": "db",
"type": "string"
},
{
"name": "sequence",
"type": [
"null",
"string"
],
"default": null
},
{
"name": "schema",
"type": "string"
},
{
"name": "table",
"type": "string"
},
{
"name": "txId",
"type": [
"null",
"long"
],
"default": null
},
{
"name": "lsn",
"type": [
"null",
"long"
],
"default": null
},
{
"name": "xmin",
"type": [
"null",
"long"
],
"default": null
}
],
"connect.name": "io.debezium.connector.postgresql.Source"
}
},
{
"name": "op",
"type": "string"
},
{
"name": "ts_ms",
"type": [
"null",
"long"
],
"default": null
},
{
"name": "transaction",
"type": [
"null",
{
"type": "record",
"name": "ConnectDefault",
"namespace": "org.apache.pulsar.kafka.shade.io.confluent.connect.avro",
"fields": [
{
"name": "id",
"type": "string"
},
{
"name": "total_order",
"type": "long"
},
{
"name": "data_collection_order",
"type": "long"
}
]
}
],
"default": null
}
],
"connect.name": "mydbinstance.public.users.Envelope"
}
}
I think there’s nothing informative in its key. A walk around would be utilize some Pulsar Function to extract the value from the original message and push it into a second topic. Make sure that topic has a schema of AVRO or JSON format, and then use the spark connector to do following processing.
It’s a little bit redundant, but if you want the pipeline flow immediately, this would be the fastest way.Oneeb
09/25/2023, 1:29 PMallowDifferentTopicSchemas flag which lets us consume without enforcing the topic schema at all.Neng
09/25/2023, 5:48 PM