This message was deleted.
# general
s
This message was deleted.
n
The existing topic’s schema and the newly produced message schema (after your config update) are incompatible. Could you try use another new topic for the messages after updating the config?
o
@Neng which specific configuration change do you mean?
json-with-envelope: False
? Thats how I originally ran it and I would get the same error.
n
your config change is fine. it’s just the topic schema is still using the old format so new messages can not be published to it. You need to either using another topic for publishing the new messages or cleanup the existing topic and reset its schema (probably deleting the whole topic is a cleaner way)
o
@Neng I made the change and deleted the topic, when I ran it again I still got:
Copy code
Received 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 JSON
n
hmm, this looks like the topic is not cleaned completely..
o
any idea how to clean it up? I have deleted the data and log folders on pulsar local and deleted the topics as well + restarted the broker too
From what I can tell, it looks like
org.apache.kafka.json.JsonConverter
isnt having any effect
@Neng I found out that I was passing the wrong configuration in the topic name and corrected it. Now the payload in Pulsar looks correct but when I now try to read from Spark connector, I get the error:
We do not support KEY_VALUE
. Any help here?
n
Just did some investigation: https://github.com/streamnative/pulsar-spark/blob/cba448939e591e8dd233e1a29ba9df06[…]d2/src/main/scala/org/apache/spark/sql/pulsar/SchemaUtils.scala Currently the pulsar-spark connector doesn’t support Key-Value schema. So it’s not able to process the data. Could you please create an issue to request this feature in the repo? Meanwhile, I also checked how does the schema look like, it should be something as follows:
Copy code
{
  "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.
o
Thanks @Neng I have created an issue: https://github.com/streamnative/pulsar-spark/issues/164. I was able to circumvent the problem by setting the
allowDifferentTopicSchemas
flag which lets us consume without enforcing the topic schema at all.
👀 1
👍 1
n
nice walk around!
🙌 1