Validator-interceptor#

Реализует интерфейсы ProducerInterceptor и ConsumerInterceptor, позволяет валидировать сообщения в форматах JSON/XML/AVRO по схемам, указанным в конфигурационном файле.

Предусловия#

Интерсептор позволяет использовать разные схемы для разных топиков. Для форматов JSON и XML поддерживается валидация схем.

Валидаторы могут работать со следующими типами сообщений:

  • JSON:

    • String;

    • Array[Byte];

    • com.fasterxml.jackson.databind.JsonNode.

  • XML:

    • String;

    • Array[Byte];

    • javax.xml.transform.Source;

    • org.w3c.dom.Document.

  • AVRO:

    • Array[Byte].

  • AVRO_JSON:

    • Array[Byte].

Последовательность выполнения#

Подключение к Kafka клиентам#

  1. Добавить актуальную версию интерсептора в зависимости проекта.

  2. Создать конфигурационный файл с настройками валидаторов и схем в соответствии с примером.

  3. Добавить настройки интерсептора к настройкам kafka клиента в соответствии с примерами.

Подключение к Spring Kafka#

Для настройки клиента kafka в Spring необходимо выполнить следующие шаги:

1. Настройка KafkaTemplate и ProducerFactory

  @Configuration
  public static class JsonConfigConsumer {

    @Bean
    public ProducerFactory<String, String> producerFactorySettingsFromFile() {

      /**
        Указываем путь до параметров `PATH_TO_PRODUCER_VALID_PROPERTIES`, дополнительные параметры можно передать через `Map.of("nameParam, valueParam")`
      */
      final var producerProperties = loadProperties(PATH_TO_PRODUCER_VALID_PROPERTIES, Map.of());

      return new DefaultKafkaProducerFactory<>(propertiesToMapConverter(producerProperties));
    }

    @Bean
    public KafkaTemplate<String, String> kafkaTemplateSettingsFromFile(ProducerFactory<String, String> producerFactorySettingsFromFile) {
      return new KafkaTemplate<>(producerFactorySettingsFromFile);
    }

    private Properties loadProperties(String path, Map<String, String> additionalSettings) {
      try {
        return FileUtils.loadPropertiesFromClasspath(path, additionalSettings);
      } catch (IOException e) {
      throw new RuntimeException(e);
      }
    }
  }

2. Настройка Listeners и ConsumerFactory

  @Configuration
  public static class JsonConfigConsumer {

        @Bean
        public ConsumerFactory<String, String> consumerFactorySettingsFromFile() {
            /**
              Указываем путь до параметров `PATH_TO_CONSUMER_VALID_PROPERTIES`, дополнительные параметры можно передать через `Map.of("nameParam, valueParam")`
            */
            Properties consumerProperties = loadProperties(PATH_TO_CONSUMER_VALID_PROPERTIES, Map.of());
            return new DefaultKafkaConsumerFactory<>(propertiesToMapConverter(consumerProperties));
        }

        @Bean
        public KafkaListenerContainerFactory<
                ConcurrentMessageListenerContainer<String, String>> kafkaListenerContainerFactorySettingsFromFile(
                ConsumerFactory<String, String> consumerFactorySettingsFromFile,
                ErrorHandler createErrorHandler
        ) {
            ConcurrentKafkaListenerContainerFactory<String, String> factory =
                    new ConcurrentKafkaListenerContainerFactory<>();
            factory.setConsumerFactory(consumerFactorySettingsFromFile);

            return factory;
        }
  }

3. Настройка Listener

  @EnableKafka
  @Component
  public class ListenersKafkaConsumerJson {

    @KafkaListener(
            id = "id-1",
            topics = "NAME_TOPIC",
            groupId = "group-id",

            /**
              Необходимо указать имя bean фабрики с нужными параметрами.
            */
            containerFactory = "kafkaListenerContainerFactorySettingsFromFile",
    )
    private void listenerConfigurationFile(String data) {
      try {
        JsonConsumerSpringTest.RESPONSE_QUEUE.put(data);
      } catch (InterruptedException e) {
            throw new RuntimeException(e);
      }
    }
  }

Примеры конфигурации#

Пример конфигурации файла со списком путей до схем#

Для каждого топика можно задать отдельную схему, а также схему по умолчанию с именем топика "*". Имя топика, который будет использоваться как топик по умолчанию, можно задать с помощью настройки interceptor.validator.default.topic.

Если схема по умолчанию отсутствует и сообщение было отправлено/получено из топика, для которого не указана схема – сообщение считается НЕВАЛИДНЫМ. Если валидация для какого-либо топика не требуется – можно настроить валидатор без схемы с типом noop.

schemas: {
    # Валидатор для топика "topic-1" (json валидатор)
    "topic-1": {
        # Тип схемы/сообщений
        "type": "json",

        # Путь до схемы на файловой системе
        "schema": "path/to/schema.json"

        # ОПЦИОНАЛЬНО Кодировка схемы
        # "schemaEncoding": "UTF-8"
    }

    # Валидатор для топика  "topic-2" (xml валидтор)
    "topic-2": {
        # Тип схемы/сообщений
        "type": "xml"

        # Путь до схемы на файловой системе
        "schema": "file://path/to/schema.xml"

        # ОПЦИОНАЛЬНО Кодировка схемы
        # "schemaEncoding": "UTF-8"
    }

    # Валидатор для топика "topic-3" (avro валидтор)
    "topic-3": {
        # Тип схемы/сообщений
        "type": "avro"

        # Путь до схемы в classpath
        "schema": "classpath://path/to/schema.avsc"

        # ОПЦИОНАЛЬНО Кодировка схемы
        # "schemaEncoding": "UTF-8"
    }

    # Валидатор для топика "topic-4" (avro_json валидатор без использования schema registry)
    "topic-4": {
        # Тип схемы/сообщений
        # Валидация avro с помощью json-схемы
        "type": "avro_json"
        # Использовать schema registry для загрузки avro-схем
        "use.schema.registry": "false"

        # Путь до avro схемы
        "avro.schema": "classpath://path/to/schema.avsc"

        # Путь до avro схемы десериализации. Если не задано, для десериализации используется 'avro.schema'
        "deserialization.avro.schema": "classpath://path/to/schema.avsc"

        # Путь до json схемы в classpath
        "schema": "classpath://path/to/schema.json"

        # Режим сериализации сообщений
        "serialization.mode": "attach.schema.id"

        # ОПЦИОНАЛЬНО Всегда использовать указанный schema id (-1 отключает настройку)
        # "use.schema.id": "-1"

        # ОПЦИОНАЛЬНО Использовать указанный schema id если он отсутствует в сообщении (-1 отключает настройку)
        # "default.schema.id": "-1"

        # ОПЦИОНАЛЬНО Кодировка схемы
        # "schemaEncoding": "UTF-8"
    }

    # Валидатор для топика "topic-5" (avro_json валидатор с использованием schema registry)
    "topic-5": {
        # Тип схемы/сообщений
        # Валидация avro с помощью json-схемы
        "type": "avro_json"

        # Использовать schema registry для загрузки avro-схем
        "use.schema.registry": "false"

        # Url schema-registry
        "schema.registry.url": "http://localhost:8081"

        # Настройки ssl для schema-registry
        # Данные настройки можно указать в настройках kafka-клиента (подробнее ниже в примере конфигурации Kafka Consumer/Producer)
        # Настройки ssl в этом файле имеют больший приоритет, чем настройки kafka-клиента
        # Настройки аналогичны настройкам ssl kafka-клиента
        "schema.registry.ssl.keystore.location": "path/to/keystore.jks"
        "schema.registry.ssl.keystore.password": "password"
        "schema.registry.ssl.truststore.location": "path/to/truststore.jks"
        "schema.registry.ssl.truststore.password": "password"
        "schema.registry.ssl.endpoint.identification.algorithm: ""

        # Путь до json схемы в classpath
        "schema": "classpath://path/to/schema.json"

        # Режим сериализации сообщений
        "serialization.mode": "attach.schema.id"

        # ОПЦИОНАЛЬНО Всегда использовать указанный schema id (-1 отключает настройку)
        # "use.schema.id": "-1"

        # ОПЦИОНАЛЬНО Использовать указанный schema id если он отсутствует в сообщении (-1 отключает настройку)
        # "default.schema.id": "-1"

        # ОПЦИОНАЛЬНО Кодировка схемы
        # "schemaEncoding": "UTF-8"
    }

    # Валидатор для всех остальных топиков
    "*": {
        # Валидатор с типом noop считает любое сообщение валидным
        type: "noop"
    }
}

Пример конфигурации Kafka Producer#

...
# 1) Подключить интерсептор
interceptor.classes=ru.sbt.ss.kafka.validator.interceptor.ValidatorProducerInterceptor

# 2) Указать путь до конфигурационного файла с настройками валидаторов

# Загрузка из файловой системы
interceptor.validator.config=path/to/validator.conf
# Загрузка из classpath
# interceptor.validator.config=classpath://path/to/validator.conf

# Для прямой загрузки avro схемы валидации 
# Загрузка схемы, указывается имя топика, знак `@` и путь до схемы, дальше знак `;` для разделения разных топиков и схем 
interceptor.validator.schemas = TOPIC@path/to/schema.avsc;TOPIC2@path/to/schema2.avsc (если используется одна схема для всех топиков interceptor.validator.schemas = path/to/schema.avsc)
# Тип схемы
interceptor.validator.type = avro
# ОПЦИОНАЛЬНО Кодировка схемы
interceptor.validator.schema.encoding = UTF-8

# 3) ОПЦИОНАЛЬНО Включить валидацию схем при запуске интерсептора
# interceptor.validator.schema.validation.enabled=true

# 4) ОПЦИОНАЛЬНО Настроить режим валидации схем (по умолчанию true)
# interceptor.validator.fail.on.invalid.schema = false

# 5) ОПЦИОНАЛЬНО Настроить режим работы продьсера при ошибках валидации (по умолчанию failOnValue)
# interceptor.validator.mode = failOnValue

# 6) ОПЦИОНАЛЬНО Включить логирование схем при загрузке
# interceptor.validator.print.schemas = true

# 7) ОПЦИОНАЛЬНО Настроить провайдер для загрузки схем
# interceptor.validator.schema.provider.class =

# 8) ОПЦИОНАЛЬНО Настроить schema resolver для выбора схемы валидации в зависимости от состава сообщения
# По умолочанию используется ru.sbt.ss.kafka.validator.interceptor.TopicSchemaResolver
# при данной настройке схема выбирается по заголовку или в зависимости от реализации интерфейса ValidatorInterceptorSchemasResolver
# interceptor.validator.schema.resolver.class = ru.sbt.ss.kafka.validator.interceptor.SchemaResolver

# 9) ОПЦИОНАЛЬНО При использовании avro_validator с schema-registry указать настройки ssl для подключения к schema-registry
# Данные настройки можно также указать в файле с настройками валидаторов,
# но при использовании параметров из конфигурации kafka-клиента можно использовать config provider'ы, например для шифрования паролей или получения сертификатов из vault
# validator.schema.registry.ssl.keystore.location = path/to/keystore.jks
# validator.schema.registry.ssl.keystore.password = password
# validator.schema.registry.ssl.truststore.location = path/to/truststore.jks
# validator.schema.registry.ssl.truststore.password = password
# validator.schema.registry.ssl.endpoint.identification.algorithm =

# 10) ОПЦИОНАЛЬНО При использовании json_validator
# Отключить проверку "uniqueItems=true" для массивов при валидации схемы
# interceptor.validator.schema.validation.disable.uniqueitems.check = true
#
# Разрешить использование $ref (по умолчанию false)
# interceptor.validator.json.schema.validation.refs.enabled = false
#
# Заменить ссылки "$ref" на объекты по ссылкам при загрузке схемы (по умолчанию значение настройки interceptor.validator.json.schema.validation.refs.enabled)
# Используется для определения цикличных и рекурсивных ссылок
# interceptor.validator.json.schema.dereferencing.enabled = true
#
# Указать язык, используемый в сообщениях об ошибках при валидации схем (по умолчанию en, возможно указать ru для использования русского языка)
# interceptor.validator.json.error.messages.locale = en

# 11) ОПЦИОНАЛЬНО
# Включить логирование интерсептора при подписи сообщений
# interceptor.trace.logger.enabled = true

# Задать имя логгера
# interceptor.trace.logger.name = ClassName[clientId]

# Включить публикацию jmx-метрик количества успешно и неудачно обработанных сообщений
# interceptor.jmx.metrics.enabled = true

# 12) ОПЦИОНАЛЬНО

# Включить публикацию телеметрии
# interceptor.telemetry.enabled = true

# Задать путь к файлу конфигурации телеметрии
# interceptor.telemetry.config.path = /path/to/config/file

# Включить отображение спанов валидации сообщений (по умолчанию false)
# interceptor.telemetry.callback.enabled = true

Для использования телеметрии необходимо подключить библиотеку kafka-telemetry-interceptor. Подробнее в разделе Подключение телеметрии.

Пример конфигурации Kafka Consumer#

# 1) Подключить интерсептор
interceptor.classes=ru.sbt.ss.kafka.validator.interceptor.ValidatorConsumerInterceptor

# 2) Указать путь до конфигурационного файла с настройками валидаторов
# Загрузка из файловой системы
interceptor.validator.config=path/to/validator.conf
# Загрузка из classpath
# interceptor.validator.config=classpath://path/to/validator.conf

# 3) ОПЦИОНАЛЬНО Включить валидацию схем при запуске интерсептора
# interceptor.validator.schema.validation.enabled=true

# 4) ОПЦИОНАЛЬНО Настроить режим валидации схем (по умолчанию true)
# interceptor.validator.fail.on.invalid.schema = false

# 5) ОПЦИОНАЛЬНО Настроить режим работы консюмера при ошибках валидации (по умолчанию failOnValue)
# interceptor.validator.mode = failOnValue

# 6) ОПЦИОНАЛЬНО Включить логирование схем при загрузке
# interceptor.validator.print.schemas = true

# 7) ОПЦИОНАЛЬНО Настроить провайдер для загрузки схем
# interceptor.validator.schema.provider.class =

# 8) ОПЦИОНАЛЬНО Настроить schema resolver для выбора схемы валидации в зависимости от состава сообщения
# По умолочанию используется ru.sbt.ss.kafka.validator.interceptor.TopicSchemaResolver
# при данной настройке схема выбирается по заголовку или в зависимости от реализации интерфейса ValidatorInterceptorSchemasResolver
# interceptor.validator.schema.resolver.class = ru.sbt.ss.kafka.validator.interceptor.SchemaResolver
# interceptor.validator.schema.resolver.class = ru.sbt.ss.kafka.validator.interceptor.SchemaResolver

# 9) ОПЦИОНАЛЬНО При использовании avro_validator с schema-registry указать настройки ssl для подключения к schema-registry
# Данные настройки можно также указать в файле с настройками валидаторов,
# но при использовании параметров из конфигурации kafka-клиента можно использовать config provider'ы, например для шифрования паролей или получения сертификатов из vault
# validator.schema.registry.ssl.keystore.location = path/to/keystore.jks
# validator.schema.registry.ssl.keystore.password = password
# validator.schema.registry.ssl.truststore.location = path/to/truststore.jks
# validator.schema.registry.ssl.truststore.password = password
# validator.schema.registry.ssl.endpoint.identification.algorithm =

# 10) ОПЦИОНАЛЬНО При использовании json_validator
# Отключить проверку "uniqueItems=true" для массивов при валидации схемы
# interceptor.validator.schema.validation.disable.uniqueitems.check = true
#
# Разрешить использование $ref (по умолчанию false)
# interceptor.validator.json.schema.validation.refs.enabled = false
#
# Заменить ссылки "$ref" на объекты по ссылкам при загрузке схемы (по умолчанию значение настройки interceptor.validator.json.schema.validation.refs.enabled)
# Используется для определения цикличных и рекурсивных ссылок
# interceptor.validator.json.schema.dereferencing.enabled = true
#
# Указать язык, используемый в сообщениях об ошибках при валидации схем (по умолчанию en, возможно указать ru для использования русского языка)
# interceptor.validator.json.error.messages.locale = en

# 11) ОПЦИОНАЛЬНО
# Включить логирование интерсептора при подписи сообщений
# interceptor.trace.logger.enabled = true

# Задать имя логгера
# interceptor.trace.logger.name = ClassName[clientId]

# Включить публикацию jmx-метрик количества успешно и неудачно обработанных сообщений
# interceptor.jmx.metrics.enabled = true

# 12) ОПЦИОНАЛЬНО

# Включить публикацию телеметрии
# interceptor.telemetry.enabled = true

# Задать путь к файлу конфигурации телеметрии
# interceptor.telemetry.config.path = /path/to/config/file

# Включить отображение спанов валидации сообщений (по умолчанию false)
# interceptor.telemetry.callback.enabled = true

Для использования телеметрии необходимо подключить библиотеку kafka-telemetry-interceptor. Подробнее в разделе Подключение телеметрии.

Интерфейс schema resolver для выбора схемы валидации#

Перехватчик предоставляет интерфейс ru.sbt.ss.kafka.validator.ValidatorInterceptorSchemasResolver, позволяющий использовать свою реализацию schema resolver для выбора схем в зависимости от состава сообщения.

Загрузка происходит следующим образом:

  1. Перехватчик получает имя класса провайдера из конфигурации kafka-клиента с помощью параметра interceptor.validator.schema.resolver.class.

  2. Перехватчик создает новый экземпляр провайдера с помощью конструктора по умолчанию без параметров.

  3. Перехватчик вызывает метод void configure(Map<String, String> configs) с конфигурацией, переданной kafka-клиенту.

  4. Перехватчик вызывает метод String getSchemas(ConsumerRecord<K, V> record) или String getSchemas(ProducerRecord<K, V> record) в зависимости от типа аргумента и по умолчанию возвращает имя топика.

По умолчанию используется interceptor.validator.schema.resolver.class=ru.sbt.ss.kafka.validator.interceptor.TopicSchemaResolver, при этом схема выбирается по имени топика. Для выбора схемы в зависимости от заголовка можно настроить interceptor.validator.schema.resolver.class=ru.sbt.ss.kafka.validator.interceptor.SchemaResolver.

package ru.sbt.ss.kafka.validator;

import ru.sbt.ss.validator.schema.SchemaProvider;

import java.util.Map;

/**
 * Schemas resolver for validator-interceptor.
 *
 * Implementations should have default constructor without arguments.
 *
 * Interceptor will create and configure instance of this interface,
 * then call getSchemas() method and returns topic name by default
 * or a specific ConsumerRecord/ProducerRecord parameter
 * by which the validation scheme will be selected.
 */
public interface ValidatorInterceptorSchemasResolver {

	/**
	 * This methods is called once after interceptor creation
	 *
     * @param record kafka ConsumerRecord
	 * @return topic name by default or a specific ConsumerRecord parameter
	 */
    String getSchemas(ConsumerRecord<K, V> record);

    /**
     * This methods is called once after interceptor creation
     *
     * @param record kafka ProducerRecord
     * @return topic name by default or a specific ProducerRecord parameter
     */
    String getSchemas(ProducerRecord<K, V> record);

	/**
	 * Optional configuration method, will be called after resolver creation.
	 *
	 * @param configs kafka configuration
	 */
	void configure(Map<String, String> configs) {}

}

Пример реализации интерфейса schema resolver c выбором схемы в зависимости от определенного заголовка#

package ru.sbt.ss.kafka.validator;

import ru.sbt.ss.kafka.validator.config.ValidatorInterceptorConfig.HeaderNameKey;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.producer.ProducerRecord;

public class SchemaResolver<K, V> extends ValidatorInterceptorSchemasResolver<K, V> {

    private String header;

    @Override
    public String getSchemas(ConsumerRecord<K, V> record) {
        return new String(record.headers().lastHeader(header).value());
    }

    @Override
    public String getSchemas(ProducerRecord<K, V> record) {
        return new String(record.headers().lastHeader(header).value());
    }

    @Override
    public void configure(Map<String, String> configs) {
        header = configs.get(HeaderNameKey);
    }
}

Валидация AVRO с помощью JSON-схемы#

Валидатор с типом avro_json позволяет валидировать сообщения в формате AVRO с помощью JSON-схемы. Поддерживается два типа AVRO-сообщений:

  • стандартный AVRO-формат;

  • формат io.confluent.kafka.serializers.KafkaAvroDeserializer (confluent формат), содержит идентификатор AVRO-схемы (schemaId) в schema registry.

Валидация происходит в три этапа:

1. Десериализация сообщения в GenericRecord с помощью AVRO-схемы

Схема может быть загружена из файла или schema registry в зависимости от параметра use.schema.registry. При использовании schema registry для десериализации используется io.confluent.kafka.serializers.KafkaAvroDeserializer.

При загрузке схемы из файла можно использовать отдельную схему для десериализации сообщения deserialization.avro.schema – в этом случае необходимо указывать путь до схемы сериализации avro.schema.

Если не задан путь до отдельной схемы десереализации, то при сериализации и десериализации используется одна схема, заданная параметром avro.schema.

2. Конвертация десериализованного сообщения в JSON-формат и его валидация с помощью JSON-схемы

Механизм валидации аналогичен валидатору с типом JSON.

3. Сериализация сообщения

При сериализации опционально можно добавить schema id или схему целиком, за режим сериализации отвечает настройка serialization.mode:

  • raw – не добавлять информацию о схеме (стандартный AVRO-формат);

  • attach.schema.id – добавлять в сообщение schema registry schemaId (сообщение конвертируется в формат io.confluent.kafka.serializers.KafkaAvroDeserializer);

  • attach.schema – добавлять в сообщение схему целиком. В этом случае для сериализации используется org.apache.avro.file.DataFileWriter.DataFileWriter.

Для десериализации такого сообщения можно использовать org.apache.avro.file.DataFileReader.DataFileReader:

GenericDatumReader<GenericRecord> datumReader = new GenericDatumReader<>();
DataFileReader<GenericRecord> dataFileReader = new DataFileReader<>(new SeekableByteArrayInput(serializedAvroBytes), datumReader);
GenericRecord avroRecord = dataFileReader.next();

Пример валидации AVRO-типов с помощью JSON-схемы#

AVRO-схема:

{
  "type": "record",
  "name": "record_type",
  "namespace": "my.example",
  "fields": [
    {
      "name": "null_type",
      "type": "null"
    },
    {
      "name": "boolean_type",
      "type": "boolean"
    },
    {
      "name": "int_type",
      "type": "int"
    },
    {
      "name": "long_type",
      "type": "long"
    },
    {
      "name": "float_type",
      "type": "float"
    },
    {
      "name": "double_type",
      "type": "double"
    },
    {
      "name": "bytes_type",
      "type": "bytes"
    },
    {
      "name": "string_type",
      "type": "string"
    },
    {
      "name": "enum_type",
      "type" : {
        "name": "string_enum",
        "type": "enum",
        "symbols": ["enum1", "enum2"]
      }
    },
    {
      "name": "array_type",
      "type": {
        "name": "string_array",
        "type": "array",
        "items": "string"
      }
    },
    {
      "name": "map_type",
      "type": {
        "name": "string_map",
        "type": "map",
        "values": "string"
      }
    },
    {
      "name": "union_type",
      "type": ["null", "string", "int"]
    }
  ]
}

JSON-схема:

{
    "id": "http://json-schema.org/draft-04/schema#",
    "$schema": "http://json-schema.org/draft-04/schema#",
    "description": "Core schema meta-schema",
    "definitions": {
        "schemaArray": {
            "type": "array",
            "minItems": 1,
            "items": { "$ref": "#" }
        },
        "positiveInteger": {
            "type": "integer",
            "minimum": 0
        },
        "positiveIntegerDefault0": {
            "allOf": [ { "$ref": "#/definitions/positiveInteger" }, { "default": 0 } ]
        },
        "simpleTypes": {
            "enum": [ "array", "boolean", "integer", "null", "number", "object", "string" ]
        },
        "stringArray": {
            "type": "array",
            "items": { "type": "string" },
            "minItems": 1,
            "uniqueItems": true
        },
        "validateAdditionalProperties": {
            "anyOf": [
                {
                    "properties": {
                        "additionalProperties": { "const": false }
                    }
                },
                {
                                    "$opt": {
                        "ifOptionEnabled": "STRING_ADDITIONAL_PROPERTIES_ENABLED",
                        "thenUseSchema": {
                            "properties": {
                                "additionalProperties": {
                                    "type": "object",
                                    "properties": {
                                        "type": {
                                            "const": "string"
                                        }
                                    }
                                }
                            },
                            "required": [
                                "maxProperties"
                            ]
                        },
                        "invalidWhenOptionDisabled": true
                    }
                }
            ]
        },
        "requireNoAdditionalItems": {
            "properties": {
                "additionalItems": { "const": false }
            }
        },
        "requireUniqueItems": {
            "properties": {
                "uniqueItems": { "const": true }
            }
        },
        "validateDescription": {
            "properties": {
                "description": {
                    "type": "string",
                    "maxLength": 250
                }
            }
        },
        "requireMaxLengthMaximum": {
            "properties": {
                "maxLength": {
                    "maximum": 250
                }
            }
        },
        "hasType": {
            "required": ["type"]
        },
        "isString": {
            "anyOf": [
                {
                    "properties": {
                        "type": { "const": "string" }
                    }
                },
                {
                    "properties": {
                        "type": {
                            "type": "array",
                            "contains": { "const": "string" }
                        }
                    }
                }
            ]
        },
        "isArray": {
            "anyOf": [
                {
                    "properties": {
                        "type": { "const": "array" }
                    }
                },
                {
                    "properties": {
                        "type": {
                            "type": "array",
                            "contains": { "const": "array" }
                        }
                    }
                }
            ]
        },
        "isObject": {
            "anyOf": [
                {
                    "properties": {
                        "type": { "const": "object" }
                    }
                },
                {
                    "properties": {
                        "type": {
                            "type": "array",
                            "contains": { "const": "object" }
                        }
                    }
                }
            ]
        },
        "isNumber": {
            "anyOf": [
                {
                    "properties": {
                        "type": { "const": "number" }
                    }
                },
                {
                    "properties": {
                        "type": {
                            "type": "array",
                            "contains": { "const": "number" }
                        }
                    }
                }
            ]
        },
        "isInteger": {
            "anyOf": [
                {
                    "properties": {
                        "type": { "const": "integer" }
                    }
                },
                {
                    "properties": {
                        "type": {
                            "type": "array",
                            "contains": { "const": "integer" }
                        }
                    }
                }
            ]
        },
        "isNumberOrInteger": {
            "anyOf": [
                { "$ref": "#/definitions/isNumber" },
                { "$ref": "#/definitions/isInteger" }
            ]
        }
    },
    "type": "object",
    "minProperties": 1,
    "properties": {
        "id": {
            "type": "string"
        },
        "$schema": {
            "type": "string"
        },
        "title": {
            "type": "string"
        },
        "description": {
            "type": "string"
        },
        "default": {},
        "multipleOf": {
            "type": "number",
            "minimum": 0,
            "exclusiveMinimum": true
        },
        "maximum": {
            "type": "number"
        },
        "exclusiveMaximum": {
            "type": "boolean",
            "default": false
        },
        "minimum": {
            "type": "number"
        },
        "exclusiveMinimum": {
            "type": "boolean",
            "default": false
        },
        "maxLength": { "$ref": "#/definitions/positiveInteger" },
        "minLength": { "$ref": "#/definitions/positiveIntegerDefault0" },
        "pattern": {
            "type": "string",
            "format": "regex"
        },
        "additionalItems": {
            "anyOf": [
                { "type": "boolean" },
                { "$ref": "#" }
            ],
            "default": {}
        },
        "items": {
            "anyOf": [
                { "$ref": "#" },
                { "$ref": "#/definitions/schemaArray" }
            ],
            "default": {}
        },
        "maxItems": { "$ref": "#/definitions/positiveInteger" },
        "minItems": { "$ref": "#/definitions/positiveIntegerDefault0" },
        "uniqueItems": {
            "type": "boolean",
            "default": false
        },
        "maxProperties": { "$ref": "#/definitions/positiveInteger" },
        "minProperties": { "$ref": "#/definitions/positiveIntegerDefault0" },
        "required": { "$ref": "#/definitions/stringArray" },
        "additionalProperties": {
            "anyOf": [
                { "type": "boolean" },
                { "$ref": "#" }
            ],
            "default": {}
        },
        "definitions": {
            "type": "object",
            "additionalProperties": { "$ref": "#" },
            "default": {}
        },
        "properties": {
            "type": "object",
            "additionalProperties": { "$ref": "#" },
            "default": {}
        },
        "patternProperties": {
            "type": "object",
            "additionalProperties": { "$ref": "#" },
            "default": {}
        },
        "dependencies": {
            "type": "object",
            "additionalProperties": {
                "anyOf": [
                    { "$ref": "#" },
                    { "$ref": "#/definitions/stringArray" }
                ]
            }
        },
        "enum": {
            "type": "array",
            "minItems": 1,
            "uniqueItems": true
        },
        "type": {
            "anyOf": [
                { "$ref": "#/definitions/simpleTypes" },
                {
                    "type": "array",
                    "items": { "$ref": "#/definitions/simpleTypes" },
                    "minItems": 1,
                    "uniqueItems": true
                }
            ]
        },
        "format": { "type": "string" },
        "allOf": { "$ref": "#/definitions/schemaArray" },
        "anyOf": { "$ref": "#/definitions/schemaArray" },
        "oneOf": { "$ref": "#/definitions/schemaArray" },
        "not": { "$ref": "#" }
    },
    "dependencies": {
        "exclusiveMaximum": [ "maximum" ],
        "exclusiveMinimum": [ "minimum" ]
    },
    "default": {},

    "$opt": {
        "ifOptionEnabled": "$REF_DISABLED",
        "thenUseSchema": {
            "not": { "required": ["$ref"] }
        }
    },
    "if": { "$ref": "#/definitions/hasType" },
    "then": {
    "allOf": [
        {
            "$warn": {
                "allOf": [
                    { "required": ["description"] },
                    { "$ref": "#/definitions/validateDescription" }
                ]
            }
        },

        {
            "if": { "$ref": "#/definitions/isArray" },
            "then": {
                "allOf": [
                    { "required": ["maxItems", "additionalItems"] },
                    { "$ref": "#/definitions/requireNoAdditionalItems" },
                    {
                        "$opt": {
                            "ifOptionEnabled": "UNIQUE_ITEMS_CHECK_ENABLED",
                            "thenUseSchema": {
                                "required": ["uniqueItems"]
                            }
                        }
                    },
                    {
                        "$opt": {
                            "ifOptionEnabled": "UNIQUE_ITEMS_CHECK_ENABLED",
                            "thenUseSchema": {
                                "$ref": "#/definitions/requireUniqueItems"
                            }
                        }
                    }
                ]
            }
        },

        {
            "if": { "$ref": "#/definitions/isObject" },
            "then": {
                "allOf": [
                    { "required": ["additionalProperties"] },
                    { "$ref": "#/definitions/validateAdditionalProperties" }
                ]
            }
        },

        {
            "if": { "$ref": "#/definitions/isNumberOrInteger" },
            "then": {
                "allOf": [
                    { "required": ["minimum", "maximum"] }
                ]
            }
        },

        {
            "if": { "$ref": "#/definitions/isString" },
            "then": {
                "allOf": [
                    { "required": ["maxLength"] },
                    { "$warn": {"$ref": "#/definitions/requireMaxLengthMaximum" } }
                ]
            }
        }
    ]
    }
}

Настройка поведения валидации при использовании двух AVRO-JSON схем#

Для установки флагов необходимо передать списком через запятую имена (регистр не важен) в параметр interceptor.validator.compatibility.rules.allowed.

Список флагов, которые можно установить:

  • ADD_OPTIONAL_FIELD: – добавление необязательного поля со значением по умолчанию;

  • DELETE_REQUIRED_FIELD: – удаление обязательного поля;

  • DELETE_OPTIONAL_FIELD: – удаление необязательного поля;

  • BECOME_OPTIONAL: – преобразование обязательного поля в необязательное путем добавления значения по умолчанию;

  • BECOME_REQUIRED: – преобразование необязательного поля в обязательное путем удаления значения по умолчанию;

  • NONE – отключает все флаги.

Например: interceptor.validator.compatibility.rules.allowed=ADD_OPTIONAL_FIELD, delete_required_field

Поведение по умолчанию#

При пустом значении параметра interceptor.validator.compatibility.rules.allowed флаги DELETE_REQUIRED_FIELD и DELETE_OPTIONAL_FIELD включено по умолчанию.

При передаче своего списка правил в параметр interceptor.validator.compatibility.rules.allowed флаги DELETE_REQUIRED_FIELD и DELETE_OPTIONAL_FIELD будут выключены, если их не задать явно.

Для отключения всех флагов необходимо передать значение NONE (регистр не важен) в параметр interceptor.validator.compatibility.rules.allowed.

Example 1: ADD_OPTIONAL_FIELD#

Схема сериализации:

{
  "type": "record",
  "name": "User",
  "fields": [
    {"name": "name", "type": "string"}
  ]
}

Схема десериализации:

{
  "type": "record",
  "name": "User",
  "fields": [
    {"name": "name", "type": "string"},
    {"name": "age", "type": "int", "default": 0} // Новое необязательное поле со значением по умолчанию
  ]
}

Результат: Если параметр interceptor.validator.compatibility.rules.allowed не содержит флаг ADD_OPTIONAL_FIELD, то это изменение вызовет ошибку валидации.

Example 2: DELETE_REQUIRED_FIELD#

Схема сериализации:

{
  "type": "record",
  "name": "User",
  "fields": [
    {"name": "name", "type": "string"},
    {"name": "age", "type": "int"}
  ]
}

Схема десериализации:

{
  "type": "record",
  "name": "User",
  "fields": [
    {"name": "name", "type": "string"}
  ]
}

Результат: Флаг DELETE_OPTIONAL_FIELD по умолчанию включен. Схема пройдет валидацию. Если параметр сompatibility.rules.allowed задан явно и не содержит флаг DELETE_OPTIONAL_FIELD, это изменение вызовет ошибку валидации.

Поведение при ошибках#

Поведение при ошибках валидации настраивается с помощью параметра interceptor.validator.mode.

В случае, если сообщение не прошло валидацию при отправке (producer.send()), возможны следующие режимы:

  1. failOnValue (используется по умолчанию) – клиент получит сообщение-заглушку вместо невалидного сообщения, которое содержит null вместо value и выбросит исключение ru.sbt.ss.kafka.interceptors.ProducerInterceptorException при вызове record.value() (перед сериализацией). Данный способ позволяет использовать логику обработки ошибок Kafka-клиента, в том числе вызывать методы перехватчика (но не callback).

  2. failOnSend – будет выброшено исключение ru.sbt.ss.kafka.interceptors.ProducerInterceptorError (extends Throwable) c сообщением, игнорирую логику обработки ошибок Kafka-клиента.

Сообщения исключений имеют формат Error while processing record(topic: topic): *ошибка валидации в зависимости от формата сообщения*.

В случае, если сообщение не прошло валидацию при получении (consumer.poll()), поведение настраивается параметром interceptor.validator.mode:

  1. failOnValue (используется по умолчанию) – клиент получит сообщение-заглушку вместо невалидного сообщения, которое содержит null вместо value и выбросит исключение ru.sbt.ss.kafka.interceptors.ConsumerInterceptorException при вызове record.value().

  2. filter – сообщение об ошибке будет залогировано в error, клиент не получит невалидное сообщение.

  3. failOnConsume – метод consumer.poll() выбросит исключение ru.sbt.ss.kafka.interceptors.ConsumerInterceptorError, клиент не получит ни одного сообщения из пачки.

  4. addErrorHeader – в сообщение будет добавлен заголовок с сообщением об ошибке валидации (по умолчанию __interceptor.error, настраивается с помощью параметра interceptor.validator.error.header.name).

Поведение при ошибках Spring Kafka#

Поведение при ошибках валидации настраивается с помощью параметра interceptor.validator.mode. Параметр задается при создании ConsumerFactory / ProducerFactory.

Поведение при отправке:

В случае, если сообщение не прошло валидацию при отправке (KafkaTemplate.send()), возможны следующие режимы:

  1. failOnValue (используется по умолчанию) – при попытке отправить сообщение (KafkaTemplate.send()) будет выброшено исключение ru.sbt.ss.kafka.interceptors.ProducerInterceptorException. Данный способ позволяет использовать логику обработки ошибок Kafka-клиента, в том числе вызывать методы перехватчика (но не callback).

  2. failOnSend – при попытке отправить сообщение (KafkaTemplate.send()) будет выброшено исключение ru.sbt.ss.kafka.interceptors.ProducerInterceptorError (extends Throwable) c сообщением, игнорируя логику обработки ошибок Kafka-клиента.

Поведение при получении:

В случае, если сообщение не прошло валидацию при получении, поведение настраивается параметром interceptor.validator.mode:

  • Spring автоматически при помощи консьюмера отслеживает новые сообщения в топиках и передает их в Listener;

  • Валидация происходит после того, как Spring вызывает метод consumer.poll().

Возникшие исключения при валидации полученных сообщений попадают в ErrorHandler. Чтобы задать логику обработки ошибок, необходимо при создании KafkaListenerContainerFactory задать собственную реализацию ErorHandler, затем передать в метод setErrorHandler(errorHandler).

Пример реализации ErrorHandler:

new ErrorHandler() {
  @Override
  public void handle(Exception e, ConsumerRecord<?, ?> consumerRecord) {
     throw e;
  }
};
  1. failOnValue (используется по умолчанию) – будет выброшено и обработано в ErrorHandler исключение ru.sbt.ss.kafka.interceptors.ConsumerInterceptorException.

  2. filter – клиент не получит невалидное сообщение.

  3. failOnConsume – будет выброшено ru.sbt.ss.kafka.interceptors.ConsumerInterceptorError, в данной конфигурации исключение не будет обработано ErrorHandler.

  4. addErrorHeader – в сообщение будет добавлен заголовок с сообщением об ошибке валидации (по умолчанию __interceptor.error, настраивается с помощью параметра interceptor.validator.error.header.name).

Загрузка конфигурации/схем из classpath#

Для загрузки конфигурации/схем из classpath необходимо указать протокол classpath:// в пути до файла, например:

  • для файла конфигурации: interceptor.validator.config=classpath://path/to/validator.conf

  • для файла со схемой: schema: "classpath://path/to/schema.xsd"

Если XSD-схема была загружена из classpath – дополнительные схемы, указанные с помощью <include schemaLocation="additional.xsd">, тоже будут загружены из classpath, указывать префикс внутри схемы необязательно.

Относительные пути в конфигурации#

Относительные пути до файлов в конфигурации обычно разрешаются относительно директории запуска приложения.

В некоторых случаях (для автоматических установок, например для установок с помощью скриптов) может быть полезно явно указать root-директорию, относительно которой будут разрешаться относительные пути:

  • interceptor.config.root.dir (общая для всех перехватчиков с подобной настройкой);

  • interceptor.validator.config.root.dir (переопределяет общую настройку).

Данная настройка не применяется к путям в classpath (classpath://path/to/file.txt)

Пример конфигурации:

interceptor.validator.config.root.dir=/full/path/to
# interceptor.config.root.dir=/full/path/to

# /full/path/to/validator.conf
interceptor.validator.config=validator.conf

Также работает и для путей до схем в файле конфигурации валидатора:

schemas: {
 "topic-1": {
    type: "json",

    # /full/path/to/schema.conf
    schema: "schema.json"
 }

 "topic-2": {
    type: "xml"

    # /full/path/to/schema.xml
    schema: "schema.xml"
 }

Особенности импорта AVRO-схем#

Валидатор поддерживает функцию импорта схем из поля «imports», где перечисляются относительные пути до импортируемых схем от текущей схемы.

Если импортируемая схема была не найдена, то вернется ошибка вида Import <полный путь до схемы> load error in schema <полное имя схемы>: Failed to load <resource или file> '<запрашиваемый путь>' from path '<искомый путь>'.

Если в «imports» встретится циклическая ссылка на родительскую схему, то вернется ошибка вида Circle reference found in schema <полное имя схемы где найдена циклическая ссылка> with import <полный путь до схемы> and schema <полное имя схемы>.

{
  "type": "record",
  "name": "ValidThree",
  "doc": "Тестовая схема",
  "namespace": "ru.sbt.test.valid",
  "imports": [
    "ValidTwo.avsc",
    "four/ValidFour.avsc"
  ],
  "fields": [
  {
    "name": "two",
    "doc": "Два",
    "type": "ru.sbt.test.valid.ValidTwo"
  },
  {
    "name": "four",
    "doc": "Три",
    "type": "ru.sbt.test.valid.ValidFour"
  }
]
}

Аудит событий валидации#

Реализован с помощью библиотеки:

<dependency>
  <groupId>ru.sbt.ss</groupId>
  <artifactId>validator-interceptor-audit-callback_2.13</artifactId>
</dependency>

Описание библиотеки и примеры конфигурации приведены в разделе Validator-interceptor-audit-callback.

Интерфейс провайдера для загрузки схем#

Перехватчик предоставляет интерфейс ru.sbt.ss.kafka.validator.ValidatorInterceptorSchemasProvider, позволяющий использовать свою реализацию провайдера для загрузки схем валидации.

Загрузка происходит следующим образом:

  1. Перехватчик получает имя класса провайдера из конфигурации Kafka-клиента с помощью параметра intercerptor.validator.schema.provider.class.

  2. Перехватчик создает новый экземпляр провайдера c помощью конструктора по умолчанию без параметров.

  3. Перехватчик вызывает метод void configure(Map<String, String> configs) с конфигурацией, переданной Kafka-клиенту.

  4. Перехватчик вызывает метод Map<String, SchemaProvider> getSchemas() и создает валидатор для каждого топика, присутствующего в Map<String, SchemaProvider> (ключ = имя топика).

package ru.sbt.ss.kafka.validator;

import ru.sbt.ss.validator.schema.SchemaProvider;

import java.util.Map;

/**
 * Schemas provider for validator-interceptor.
 *
 * Implementations should have default constructor without arguments.
 *
 * Interceptor will create and configure instance of this interface,
 * then call getSchemas() method and initialize validators for each entry (topic) of the schemas map.
 * Schemas map should include default topic as well.
 */
public interface ValidatorInterceptorSchemasProvider {

	/**
	 * This method is called once after interceptor creation
	 *
	 * @return map of topic -> {@link SchemaProvider}
	 */
	Map<String, SchemaProvider> getSchemas();

	/**
	 * Optional configuration method, will be called after provider creation.
	 *
	 * @param configs kafka configuration
	 */
	default void configure(Map<String, String> configs) {
		// NOOP
	}

}
/**
 * Validation schema provider for {@link ru.sbt.ss.validator.Validator}
 */
public interface SchemaProvider {

	/**
	 * Schema name, used in logging/metrics
	 */
	String getName();

	/**
	 * Schema as string
	 */
	String getSchema();

	/**
	 * Schema type
	 */
	SchemaType getType();
}

Пример реализации интерфейса провайдера схем:#

package ru.sbt.ss.kafka.validator;

import ru.sbt.ss.kafka.validator.config.ValidatorInterceptorConfig;
import ru.sbt.ss.validator.schema.SchemaProvider;
import ru.sbt.ss.validator.schema.SchemaType;

import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.util.HashMap;
import java.util.Map;

public class CustomInterceptorSchemasProvider implements ValidatorInterceptorSchemasProvider {

	public static final String NOOP_VALIDATOR_ENABLED = ValidatorInterceptorConfig.ValidatorInterceptorPrefix() + "noop.validator.enabled";

	private static final SchemaProvider AVRO_SCHEMA_PROVIDER = new ClasspathSchemaProvider("avro/correct_avro_schema.avsc", SchemaType.AVRO);
	private static final SchemaProvider JSON_SCHEMA_PROVIDER = new ClasspathSchemaProvider("json/schema_first.json", SchemaType.JSON);
	private static final SchemaProvider XML_SCHEMA_PROVIDER = new ClasspathSchemaProvider("xml/schema/xml_schema.xsd", SchemaType.XML);
	private static final SchemaProvider NOOP_SCHEMA_PROVIDER = new SchemaProvider() {
		@Override
		public String getName() {
			return "noop";
		}

		@Override
		public String getSchema() {
			return null;
		}

		@Override
		public SchemaType getType() {
			return SchemaType.NOOP;
		}
	};

	private final Map<String, SchemaProvider> schemas = new HashMap<>();

	@Override
	public void configure(Map<String, String> configs) {
		schemas.put("avro-topic", AVRO_SCHEMA_PROVIDER);
		schemas.put("json-topic", JSON_SCHEMA_PROVIDER);
		schemas.put("xml-topic", XML_SCHEMA_PROVIDER);

		final boolean noopValidatorEnabled = Boolean.parseBoolean(configs.getOrDefault(NOOP_VALIDATOR_ENABLED, "true"));
		if (noopValidatorEnabled) schemas.put(ValidatorInterceptorConfig.DefaultTopicDefault(), NOOP_SCHEMA_PROVIDER);
	}

	@Override
	public Map<String, SchemaProvider> getSchemas() {
		return schemas;
	}

	public static class ClasspathSchemaProvider implements SchemaProvider {

		private final String name;

		private String schema;

		private final SchemaType type;

		public ClasspathSchemaProvider(String schemaPath, SchemaType type) {
			this.name = schemaPath;
			try {
				this.schema = new String(getClass().getClassLoader().getResourceAsStream(schemaPath).readAllBytes(), StandardCharsets.UTF_8);
			} catch (IOException e) {
				this.schema = "";
				e.printStackTrace();
			}
			this.type = type;
		}

		@Override
		public String getName() {
			return name;
		}

		@Override
		public String getSchema() {
			return schema;
		}

		@Override
		public SchemaType getType() {
			return type;
		}
	}
}

Валидация схем#

Настроить режим валидации схем можно с помощью параметра interceptor.validator.fail.on.invalid.schema. По умолчанию interceptor.validator.fail.on.invalid.schema = true и при ошибках валидации будут выброшены исключения. При настройке interceptor.validator.fail.on.invalid.schema = false – ошибки валидации логируются в ERROR и в WARN логируется сообщение об использовании невалидной схемы.

JSON#

Механизм валидации аналогичен валидации обычных сообщений, в качестве схемы используется JSON-метасхема версии draft-04. Метасхема находится в /resources, при загрузке проверяется ее hash.

Метасхема дополнена следующими ограничениями:

  1. Для любого объекта схемы обязательно хотя бы одно поле (ограничивает использование пустых объектов "{}" в схеме).

  2. По умолчанию отключено использование ссылок "$ref". Ограничение можно снять настройкой interceptor.validator.json.schema.validation.refs.enabled = true.

  3. При использовании ссылок "$ref" ограничено использование циклических и рекурсивных ссылок. Проверка выполняется с помощью полной замены ссылок $ref в схеме на объекты по ссылке. Ограничение можно снять настройкой interceptor.validator.json.schema.dereferencing.enabled = false.

  4. При использовании ссылок "$ref" разрешено использование только абсолютных ссылок на текущий документ ({"$ref": "#/path/to/schema"}). Проверка выполняется с помощью полной замены ссылок $ref в схеме на объекты по ссылке. Ограничение можно снять настройкой interceptor.validator.json.schema.dereferencing.enabled = false.

  5. Для массивов ("type": "array") обязательны поля: maxItems, additionalItems=false, uniqueItems=true. Ограничение uniqueItems=true можно снять настройкой interceptor.validator.schema.validation.disable.uniqueitems.check = true.

  6. Для объектов ("type": "object") обязательно поле additionalProperties = false. Поле также может принимать значение {"type":"string"} с указанием maxProperties, если включена настройка interceptor.validator.schema.validation.allow.string.additional.properties=true.

  7. Для чисел ("type": "number" или "type": "integer") обязательны поля minimum и maximum.

  8. Для строк ("type": "string") обязательно поле maxLength.

  9. При указании нескольких типов для поля ("type": ["object", "array"]) проверяются ограничения для всех присутствующих типов.

При неуспешной валидации будет выброшено исключение "Schema '<Имя схемы>' doesnt match meta schema:" + список ошибок валидации.

Также метасхема дополнена следующими некритичными ограничениями, не влияющими на прохождение валидации:

  1. Для любого объекта, имеющего тип ("type"), требуется аннотация "description".

  2. Максимальная длина строк ("type": "string") не должна быть больше 250 символов ("maxLength" <= 250).

Некритичные ограничения не влияют на прохождение схемой валидации, ошибки будут залогированы в WARN.

XSD#

XSD-схема будет проверена на соответствие следующим ограничениям:

  1. Для представления числовой информации нужно использовать ограничение “totalDigits” (xs:decimal и все типы, производные от него: xs:integer, xs:negativeInteger, xs:nonNegativeInteger,xs:nonPositiveInteger, xs:positiveInteger``xs:byte, xs:long, xs:int, xs:short, xs:unsignedLong, xs:unsignedInt, xs:unsignedShort, xs:unsignedByte).

  2. Запрещено использовать «any» для описания элементов (xs:any).

  3. Не допускается неограниченная длина элементов схемы (<xs:element maxOccurs="unbounded"/>).

При неуспешной валидации будет выброшено исключение: "XML Schema '$schema' validation failed: *список ошибок валидации*"

Также будут проверены следующие некритичные ограничения, не влияющие на прохождение валидации:

  1. Все элементы схем должны быть аннотированы.

  2. Максимальная длина строк не должна быть больше 250 символов.

Некритичные ограничения не влияют на прохождение схемой валидации, ошибки будут залогированы в WARN.

AVRO#

Валидация AVRO-схем не поддерживается. При использовании avro_json валидатора JSON-схема проходит валидацию аналогично JSON-валидатору.

JMX метрики#

При успешной загрузке перехватчика в JMX будут добавлены метрики с именем/типом схемы для каждого топика и флагом прохождения схемой валидации по пути в аттрибуте Value:

kafka.consumer:type=consumer-interceptor-metrics,client-id=<client-id>,interceptor=ValidatorInterceptor,topic=<topic-name>,name=schemaName
kafka.consumer:type=consumer-interceptor-metrics,client-id=<client-id>,interceptor=ValidatorInterceptor,topic=<topic-name>,name=schemaType
kafka.consumer:type=consumer-interceptor-metrics,client-id=<client-id>,interceptor=ValidatorInterceptor,topic=<topic-name>,name=validated

или

kafka.producer:type=producer-interceptor-metrics,client-id=<client-id>,interceptor=ValidatorInterceptor,topic=<topic-name>,name=schemaName
kafka.producer:type=producer-interceptor-metrics,client-id=<client-id>,interceptor=ValidatorInterceptor,topic=<topic-name>,name=schemaType
kafka.producer:type=producer-interceptor-metrics,client-id=<client-id>,interceptor=ValidatorInterceptor,topic=<topic-name>,name=validated

Для схемы по умолчанию в имени метрики будет отсутствовать имя топика topic=<topic-name>.

Значения флага validated:

  • true – валидация схем включена и успешно пройдена;

  • false – валидация схем выключена.

client-id берется из конфигурации консьюмера client.id, при отсутствии в конфигурации генерируется автоматически. При включении публикации JMX-метрик с помощью настройки interceptor.jmx.metrics.enabled = true в JMX будут добавлены метрики с количеством успешно и ошибочно обработанных сообщений и информацией о подключенном перехватчике.

kafka.producer:type=producer-interceptor-metrics,client-id=<client-id>,interceptor=ValidatorInterceptor,name=FailedProcessedMessage
kafka.producer:type=producer-interceptor-metrics,client-id=<client-id>,interceptor=ValidatorInterceptor,name=SuccessfulProcessedMessage
kafka.producer:type=producer-interceptor-metrics,client-id=<client-id>,interceptor=ValidatorInterceptor,name=Info
kafka.consumer:type=consumer-interceptor-metrics,client-id=<client-id>,interceptor=ValidatorInterceptor,name=FailedProcessedMessage
kafka.consumer:type=consumer-interceptor-metrics,client-id=<client-id>,interceptor=ValidatorInterceptor,name=SuccessfulProcessedMessage
kafka.consumer:type=consumer-interceptor-metrics,client-id=<client-id>,interceptor=ValidatorInterceptor,name=Info

Подключение к SEDR (Сервис междоменной репликации событий)#

Подключение interceptor Validator-interceptor описано в документации к компоненту SEDR в Руководстве пользователя в разделе Установка validator-interceptor к SEDR.

Результат#

Выполнено подключение перехватчика Validator-interceptor.