Примеры ошибок при отправке невалидных данных в топик со схемой#

Если для топика создана схема и включена валидация данных на уровне брокера, то при записи/чтении по умолчанию ожидаются данные в формате, совместимом с Confluent (confluent.compatible.serialization=true).

Формат Confluent ожидает в сообщении:

  • магический байт: 0x0;

  • массив байт, полученный из идентификатора схемы;

  • сериализованные данные payload.

В случае, если Производитель пытается записать данные «как есть» (без магического байта и массива байт из id схемы), то на брокере возникнет ошибка сериализации и данные не будут записаны в топик. В server.log будет добавлена ошибка:

[2023-12-15 12:58:02,170] ERROR [ReplicaManager broker=0] Error processing append operation on partition json-test-0 (kafka.server.ReplicaManager:76)
org.apache.kafka.common.errors.PolicyViolationException: Unknown magic byte!

Внимание

Eсли Производитель не использует специального сериализатора (например ru.sbrf.kafka.schemaregistry.schemas.SchemaRegistrySerde, io.confluent.kafka.serializers.KafkaAvroSerializer, io.confluent.kafka.serializers.json.KafkaJsonSchemaSerializer), то на стороне брокера в server.properties эта опция должна быть явным образом выключена для records policy:

records.policy.https-policy.confluent.compatible.serialization=false

Примеры для формата JSON#

Пример схемы:

{
  "$schema": "http://json-schema.org/draft-07/schema#",
  "title": "Value1_checks",
  "type": "object",
  "additionalProperties": false,
  "properties": {
    "i": {
      "type": "integer",
      "minimum": 0,
      "maximum": 1000
    },
    "b": {
      "type": "boolean"
    },
    "s": {
      "type": "string",
      "maxLength": 10
    },
    "inners": {
      "$ref": "#/definitions/InnerType"
    }
  },
  "required": [
    "i",
    "b"
  ],
  "definitions": {
    "InnerType": {
      "type": "object",
      "additionalProperties": false,
      "properties": {
        "i": {
          "type": "integer",
          "minimum": 0,
          "maximum": 1000
        },
        "map": {
          "type": "object",
          "additionalProperties": {
            "type": "string",
            "maxLength": 100
          }
        },
        "v": {
          "type": "string",
          "enum": ["V1", "V2", "V3"]
        }
      },
      "required": ["i"]
    }
  }
}

Пример сообщения:

{"i" : 1001, "b" : true, "s" : "string"}

Ошибка, так как для поля i установлено максимальное значение 1000:

[2023-12-15 14:40:25,249] ERROR [ReplicaManager broker=0] Error processing append operation on partition json-test-0 (kafka.server.ReplicaManager:76)
org.apache.kafka.common.errors.PolicyViolationException: #/i: 1001 is not less or equal to 1000

Пример сообщения:

{"i" : 1000, "b" : true, "s" : "string_field"}

Ошибка, так как для поля s установлено максимальное значение длины 10:

[2023-12-15 14:40:25,249] ERROR [ReplicaManager broker=0] Error processing append operation on partition json-test-0 (kafka.server.ReplicaManager:76)
org.apache.kafka.common.errors.PolicyViolationException: #/s: expected maxLength: 10, actual: 12

Пример сообщения:

{"i" : 1000, "b" : "not_boolean", "s" : "string_field"}

Ошибка, так как для поля b установлен тип boolean:

[2023-12-15 14:40:25,249] ERROR [ReplicaManager broker=0] Error processing append operation on partition json-test-0 (kafka.server.ReplicaManager:76)
org.apache.kafka.common.errors.PolicyViolationException: #/b: expected type: Boolean, found: String

Пример сообщения:

{"i" : 1000, "b" : "not_boolean", "s" : null}

Ошибка, так как для поля s установлен тип string. Тип null это отдельный тип данных в JSON-схеме:

[2023-12-15 14:40:25,249] ERROR [ReplicaManager broker=0] Error processing append operation on partition json-test-0 (kafka.server.ReplicaManager:76)
org.apache.kafka.common.errors.PolicyViolationException: #/s: expected type: String, found: Null

Если требуется передавать помимо типа string еще null, укажите в схеме:

"s" : {"type" : ["string", "null"]}

Пример сообщения:

{"i" : 1000}

Ошибка, так как не передано обязательное поле b:

[2023-12-15 14:40:25,249] ERROR [ReplicaManager broker=0] Error processing append operation on partition json-test-0 (kafka.server.ReplicaManager:76)
org.apache.kafka.common.errors.PolicyViolationException: #: required key [b] not found

Пример сообщения:

{"i" : 1000, "b" : true, "unexpected" : "custom"}

Ошибка, так как передано необъявленное в схеме поле unexpected:

[2023-12-15 14:40:25,249] ERROR [ReplicaManager broker=0] Error processing append operation on partition json-test-0 (kafka.server.ReplicaManager:76)
org.apache.kafka.common.errors.PolicyViolationException: #: extraneous key [unexpected] is not permitted

Примеры для формата AVRO#

Пример схемы:

{
  "type": "record",
  "name": "ru.sbrf.kafka.schemaregistry.schemas.avro.Value1",
  "fields": [
    {
      "type": "int",
      "name": "i"
    },
    {
      "type": "boolean",
      "name": "b"
    },
    {
      "type": ["null", "string"],
      "name": "s"
    },
    {
      "name": "inners",
      "type": [{
        "type": "record",
        "name": "ru.sbrf.kafka.schemaregistry.schemas.avro.InnerType",
        "fields": [
          {
            "type": "int",
            "name": "i"
          },
          {
            "type": {
              "type": "map",
              "values": "string"
            },
            "name": "map"
          },
          {
            "type": {
              "type": "enum",
              "name": "ru.sbrf.kafka.schemaregistry.schemas.avro.ValuesEnum",
              "symbols": [
                "V1",
                "V2",
                "V3"
              ]
            },
            "name": "v"
          }
        ]
      }, "null"]
    }
  ]
}

Входящие данные: AVRO объект, который не содержит все not-null поля

Возможные ошибки:

[2023-12-15 15:35:29,765] ERROR [ReplicaManager broker=0] Error processing append operation on partition avro-test-0 (kafka.server.ReplicaManager:76)
org.apache.kafka.common.errors.PolicyViolationException: java.io.EOFException
[2023-12-15 15:35:29,765] ERROR [ReplicaManager broker=0] Error processing append operation on partition avro-test-0 (kafka.server.ReplicaManager:76)
org.apache.kafka.common.errors.PolicyViolationException: Index 6 out of bounds for length 2