Перейти к содержанию

Transformations

DataFlow Operator поддерживает трансформации сообщений: для каждого сообщения цепочка из transformations применяется по порядку списка (трансформер 0..N). Между сообщениями процессор может параллелить transforms при spec.transformWorkers > 1; порядок записи в sink сохраняется. См. Архитектура — concurrency пайплайна. Для доступа к полям используется gjson JSONPath.

Обзор трансформаций

Трансформация Описание Вход Выход
Timestamp Добавляет поле с временной меткой 1 сообщение 1 сообщение
Flatten Разворачивает массив в отдельные сообщения 1 сообщение N сообщений
Filter Оставляет сообщения по условию (истинность / == / != / metadata.* / && || !) 1 сообщение 0 или 1 сообщение
Mask Маскирует указанные поля 1 сообщение 1 сообщение
Router Отправляет совпадающие сообщения в альтернативные приёмники 1 сообщение 0 или 1 сообщение
Select Оставляет только указанные поля 1 сообщение 1 сообщение
Remove Удаляет указанные поля 1 сообщение 1 сообщение
SnakeCase Преобразует ключи в snake_case 1 сообщение 1 сообщение
CamelCase Преобразует ключи в CamelCase 1 сообщение 1 сообщение
DebeziumUnwrap Распаковывает Debezium envelope в строку таблицы 1 сообщение 0 или 1 сообщение
ReplaceField Переименовывает поля; опционально include/exclude без сплющивания 1 сообщение 1 сообщение
HeadersToPayload Копирует заголовки Kafka/сообщения в поля JSON payload 1 сообщение 1 сообщение
StructFlatten Сплющивает вложенные JSON-объекты в одноуровневую карту 1 сообщение 1 сообщение
ExtractField Заменяет payload значением одного поля 1 сообщение 1 сообщение
HoistField Оборачивает весь payload под одним ключом верхнего уровня 1 сообщение 1 сообщение
Cast Приводит значения полей к целевым типам (string / int64 / float64 / bool / null) 1 сообщение 1 сообщение (или skip при ошибке cast)
Timezone Конвертирует temporal-поля в целевую IANA-зону или UTC offset 1 сообщение 1 сообщение (или skip при ошибке парсинга)
InsertField Вставляет или перезаписывает поля: литералы, ${metadata.*}, ${now}, json:<raw> 1 сообщение 1 сообщение

Синтаксис условий (filter, router, when)

Общий DSL для filter.config.condition, условий маршрутов router и transformations[].when:

  • Истинность: $.field / field — true, если поле в payload существует и «истинно» (булево true, непустая строка, ненулевое число).
  • Metadata: metadata.<key> — те же проверки по Metadata сообщения (например Kafka topic, partition, offset). Путь $.metadata... — только для JSON-поля с именем metadata.
  • Равенство / неравенство: == / != с кавычками ('value' / "value") или без (true, 5).
  • Комбинаторы: && (и), || (или; && связывается сильнее), унарный !, скобки (...).

Отсутствующее поле не проходит ни истинность, ни сравнение. Пустое условие — false.

Примеры: $.status == 'active' && metadata.topic == 'orders', $.payload.op == 'u' || $.payload.op == 'c', !($.level == 'debug').

Условные трансформации (when)

Опциональный предикат на любом шаге. Если when задан и равен false, шаг пропускается, сообщение идёт дальше без изменений (passthrough). В отличие от filter, false у when не отбрасывает сообщение — аналог predicate у Kafka Connect SMT.

transformations:
  - type: debeziumUnwrap
    when: "$.payload.op == 'u' || $.payload.op == 'c'"
    config:
      addOperationFields: true
  - type: mask
    when: "metadata.topic == 'dbserver1.public.orders'"
    config:
      fields: ["payload.after.email"]
      maskChar: "*"

Timestamp

Добавляет поле с временной меткой к каждому сообщению. Полезно для отслеживания времени обработки сообщений.

Конфигурация

transformations:
  - type: timestamp
    config:
      # Имя поля для временной метки (опционально, по умолчанию: created_at)
      fieldName: created_at

      # Формат временной метки (опционально, по умолчанию: RFC3339)
      # Поддерживаются все форматы Go time package
      format: RFC3339

Формат

Поле format задаёт формат времени Go. По умолчанию — RFC3339 (например, 2006-01-02T15:04:05Z07:00). Можно указать RFC3339, RFC3339Nano или свой шаблон, например 2006-01-02 15:04:05.

Примеры

Базовое использование

transformations:
  - type: timestamp
    config:
      fieldName: processed_at

Входное сообщение:

{
  "id": 1,
  "name": "Test"
}

Выходное сообщение:

{
  "id": 1,
  "name": "Test",
  "processed_at": "2024-01-15T10:30:00Z"
}

Кастомный формат

transformations:
  - type: timestamp
    config:
      fieldName: timestamp
      format: "2006-01-02 15:04:05"

Выходное сообщение:

{
  "id": 1,
  "timestamp": "2024-01-15 10:30:00"
}

Unix timestamp

transformations:
  - type: timestamp
    config:
      fieldName: unix_time
      format: Unix

Выходное сообщение:

{
  "id": 1,
  "unix_time": "1705312200"
}

Flatten

Разворачивает массив в отдельные сообщения, сохраняя остальные поля исходного сообщения. Каждый элемент массива подставляется в корень; объекты разворачиваются в ключи верхнего уровня. Если поле не массив — сообщение возвращается без изменений. Поддерживаются массивы в формате Avro (объект с ключом array).

Конфигурация

transformations:
  - type: flatten
    config:
      # JSONPath к массиву для развертывания (обязательно)
      field: items

Примеры

Простое развертывание

transformations:
  - type: flatten
    config:
      field: items

Входное сообщение:

{
  "order_id": 12345,
  "customer": "John Doe",
  "items": [
    {"product": "Apple", "quantity": 5},
    {"product": "Banana", "quantity": 3}
  ]
}

Выходные сообщения:

{
  "order_id": 12345,
  "customer": "John Doe",
  "product": "Apple",
  "quantity": 5
}

{
  "order_id": 12345,
  "customer": "John Doe",
  "product": "Banana",
  "quantity": 3
}

Вложенные массивы

transformations:
  - type: flatten
    config:
      field: orders.items

Входное сообщение:

{
  "customer_id": 100,
  "orders": {
    "items": [
      {"sku": "SKU001", "price": 10.99},
      {"sku": "SKU002", "price": 5.99}
    ]
  }
}

Выходные сообщения:

{
  "customer_id": 100,
  "orders": {},
  "sku": "SKU001",
  "price": 10.99
}

Комбинация с другими трансформациями

transformations:
  - type: flatten
    config:
      field: rowsStock
  - type: timestamp
    config:
      fieldName: created_at

Это создаст отдельное сообщение для каждого элемента массива, каждое с добавленной временной меткой.

Filter

Оставляет только сообщения, которые удовлетворяют условию. Остальные отбрасываются.

Синтаксис условия

См. Синтаксис условий (тот же DSL, что у router и when).

Конфигурация

transformations:
  - type: filter
    config:
      # JSONPath (истинность) или сравнение == / != (обязательно)
      condition: "$.status == 'active'"

JSONPath

Используется библиотека gjson: $.field, $.nested.field, $.array[0] и т.д. Префикс $. необязателен.

Примеры

Простая фильтрация (истинность)

transformations:
  - type: filter
    config:
      condition: "$.active"

Входные сообщения:

{"id": 1, "active": true}   // ✅ Проходит
{"id": 2, "active": false}  // ❌ Отфильтровывается
{"id": 3}                   // ❌ Отфильтровывается (нет поля)

Фильтрация по равенству

transformations:
  - type: filter
    config:
      condition: "$.status == 'active'"

Входные сообщения:

{"status": "active"}    // ✅ Проходит
{"status": "inactive"}  // ❌ Отфильтровывается
{"status": ""}          // ❌ Отфильтровывается
{}                      // ❌ Отфильтровывается (нет поля)

Фильтрация по неравенству

transformations:
  - type: filter
    config:
      condition: "$.status != 'deleted'"

Входные сообщения:

{"status": "active"}   // ✅ Проходит
{"status": "deleted"}  // ❌ Отфильтровывается
{}                     // ❌ Отфильтровывается (нет поля)

Фильтрация по вложенному пути

transformations:
  - type: filter
    config:
      condition: "$.user.status == 'active'"

Входное сообщение:

{
  "user": {
    "status": "active"
  }
}

Результат: Сообщение проходит, если user.status равен active.

Mask

Маскирует чувствительные данные в указанных полях. Поддерживает сохранение длины или полное замещение символами.

Конфигурация

transformations:
  - type: mask
    config:
      # Список JSONPath выражений к полям для маскирования (обязательно)
      fields:
        - password
        - email
        - creditCard

      # Символ для маскирования (опционально, по умолчанию: *)
      maskChar: "*"

      # Сохранять оригинальную длину (опционально, по умолчанию: false)
      keepLength: true

Примеры

Маскирование с сохранением длины

transformations:
  - type: mask
    config:
      fields:
        - password
        - email
      keepLength: true

Входное сообщение:

{
  "id": 1,
  "username": "john",
  "password": "secret123",
  "email": "john@example.com"
}

Выходное сообщение:

{
  "id": 1,
  "username": "john",
  "password": "*********",
  "email": "****************"
}

Маскирование с фиксированной длиной

transformations:
  - type: mask
    config:
      fields:
        - password
      keepLength: false
      maskChar: "X"

Входное сообщение:

{
  "password": "verylongpassword123"
}

Выходное сообщение:

{
  "password": "XXX"
}

Маскирование вложенных полей

transformations:
  - type: mask
    config:
      fields:
        - user.password
        - payment.cardNumber
      keepLength: true

Входное сообщение:

{
  "user": {
    "password": "secret"
  },
  "payment": {
    "cardNumber": "1234567890123456"
  }
}

Выходное сообщение:

{
  "user": {
    "password": "******"
  },
  "payment": {
    "cardNumber": "****************"
  }
}

Router

Маршрутизирует сообщения в разные приёмники по условиям. Первое совпавшее условие задаёт приёмник; если ни одно не подошло — используется основной sink.

Синтаксис условий

См. Синтаксис условий (тот же DSL, что у filter и when).

Условия проверяются по порядку; первое совпадение определяет маршрут.

Конфигурация

transformations:
  - type: router
    config:
      routes:
        - condition: "$.level == 'error'"
          sink:
            type: kafka
            config:
              brokers: ["localhost:9092"]
              topic: error-topic
        - condition: "$.level == 'warning'"
          sink:
            type: postgresql
            config:
              connectionString: "..."
              table: warnings

Особенности

  • Условия проверяются в порядке указания; первое совпадение определяет приёмник
  • Если ни одно условие не совпало, сообщение идёт в основной sink

Примеры

Маршрутизация по уровню логирования

transformations:
  - type: router
    config:
      routes:
        - condition: "$.level"
          sink:
            type: kafka
            config:
              brokers: ["localhost:9092"]
              topic: error-logs

Входные сообщения:

{"level": "error", "message": "Critical error"}     // → error-logs топик
{"level": "info", "message": "Info message"}       // → основной приемник
{"level": "warning", "message": "Warning"}          // → основной приемник

Множественная маршрутизация

transformations:
  - type: router
    config:
      routes:
        - condition: "$.type"
          sink:
            type: kafka
            config:
              brokers: ["localhost:9092"]
              topic: events-topic
        - condition: "$.priority"
          sink:
            type: postgresql
            config:
              connectionString: "postgres://..."
              table: high_priority_events

Входные сообщения:

{"type": "event", "data": "..."}           // → events-topic
{"priority": "high", "data": "..."}       // → high_priority_events таблица
{"data": "..."}                           // → основной приемник

Комбинация с другими трансформациями

transformations:
  - type: timestamp
    config:
      fieldName: processed_at
  - type: router
    config:
      routes:
        - condition: "$.level"
          sink:
            type: kafka
            config:
              brokers: ["localhost:9092"]
              topic: errors

Сначала добавляется временная метка, затем сообщение маршрутизируется.

Select

Оставляет только указанные поля; остальные удаляются. Для каждого поля используется JSONPath; в выводе ключом становится последний сегмент пути (например, user.name → ключ name), поэтому результат плоский.

Конфигурация

transformations:
  - type: select
    config:
      # Список JSONPath к полям, которые нужно оставить (обязательно)
      fields:
        - id
        - name
        - email

Примеры

Простой выбор полей

transformations:
  - type: select
    config:
      fields:
        - id
        - name
        - email

Входное сообщение:

{
  "id": 1,
  "name": "John Doe",
  "email": "john@example.com",
  "password": "secret",
  "internal_id": 999
}

Выходное сообщение:

{
  "id": 1,
  "name": "John Doe",
  "email": "john@example.com"
}

Выбор вложенных полей (результат — плоский)

transformations:
  - type: select
    config:
      fields:
        - user.id
        - user.name
        - metadata.timestamp

Входное сообщение:

{
  "user": {
    "id": 1,
    "name": "John",
    "email": "john@example.com"
  },
  "metadata": {
    "timestamp": "2024-01-15T10:30:00Z",
    "source": "api"
  }
}

Выходное сообщение (ключи — последние сегменты путей):

{
  "id": 1,
  "name": "John",
  "timestamp": "2024-01-15T10:30:00Z"
}

Remove

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

Конфигурация

transformations:
  - type: remove
    config:
      # Список JSONPath выражений к полям для удаления (обязательно)
      fields:
        - password
        - internal_id
        - secret_token

Примеры

Удаление чувствительных полей

transformations:
  - type: remove
    config:
      fields:
        - password
        - creditCard
        - ssn

Входное сообщение:

{
  "id": 1,
  "name": "John Doe",
  "password": "secret",
  "creditCard": "1234-5678-9012-3456",
  "ssn": "123-45-6789"
}

Выходное сообщение:

{
  "id": 1,
  "name": "John Doe"
}

Удаление вложенных полей

transformations:
  - type: remove
    config:
      fields:
        - user.password
        - metadata.internal

Входное сообщение:

{
  "user": {
    "id": 1,
    "name": "John",
    "password": "secret"
  },
  "metadata": {
    "timestamp": "2024-01-15",
    "internal": "secret"
  }
}

Выходное сообщение:

{
  "user": {
    "id": 1,
    "name": "John"
  },
  "metadata": {
    "timestamp": "2024-01-15"
  }
}

Порядок применения

Трансформации применяются последовательно в порядке, указанном в списке transformations. Каждая трансформация получает результат предыдущей.

Пример последовательности

transformations:
  # 1. Развернуть массив
  - type: flatten
    config:
      field: items

  # 2. Добавить временную метку
  - type: timestamp
    config:
      fieldName: created_at

  # 3. Отфильтровать неактивные
  - type: filter
    config:
      condition: "$.active"

  # 4. Удалить внутренние поля
  - type: remove
    config:
      fields:
        - internal_id
        - debug_info

  # 5. Выбрать только нужные поля
  - type: select
    config:
      fields:
        - id
        - name
        - created_at

Рекомендации по порядку

  1. Flatten должен быть первым, если нужно развернуть массивы
  2. Filter применяйте рано, чтобы уменьшить объем обрабатываемых данных
  3. SnakeCase/CamelCase применяйте после Select/Remove, но перед отправкой в приемник
  4. Mask/Remove применяйте перед Select для безопасности
  5. Select применяйте в конце для финальной очистки
  6. Timestamp можно применять в любом месте, но обычно в начале или конце
  7. Router обычно применяется в конце, после всех других трансформаций

Комбинированные примеры

Обработка заказов

transformations:
  # Развернуть товары в отдельные сообщения
  - type: flatten
    config:
      field: items

  # Добавить временную метку
  - type: timestamp
    config:
      fieldName: processed_at

  # Фильтровать только оплаченные заказы
  - type: filter
    config:
      condition: "$.status == 'paid'"

  # Удалить чувствительные данные
  - type: remove
    config:
      fields:
        - customer.creditCard
        - customer.cvv

Обработка логов

transformations:
  # Добавить временную метку
  - type: timestamp
    config:
      fieldName: timestamp

  # Маскировать IP адреса
  - type: mask
    config:
      fields:
        - ip_address
      keepLength: true

  # Маршрутизировать ошибки
  - type: router
    config:
      routes:
        - condition: "$.level"
          sink:
            type: kafka
            config:
              brokers: ["localhost:9092"]
              topic: error-logs

Нормализация имен полей

transformations:
  # Выбрать нужные поля
  - type: select
    config:
      fields:
        - firstName
        - lastName
        - email
        - address

  # Преобразовать в snake_case для PostgreSQL
  - type: snakeCase
    config:
      deep: true

Входное сообщение:

{
  "firstName": "John",
  "lastName": "Doe",
  "email": "john@example.com",
  "address": {
    "streetName": "Main St",
    "zipCode": "12345"
  }
}

Выходное сообщение:

{
  "first_name": "John",
  "last_name": "Doe",
  "email": "john@example.com",
  "address": {
    "street_name": "Main St",
    "zip_code": "12345"
  }
}

JSONPath поддержка

Все трансформации, работающие с полями, поддерживают JSONPath синтаксис:

  • $.field - корневое поле
  • $.nested.field - вложенное поле
  • $.array[0] - элемент массива по индексу
  • $.array[*] - все элементы массива
  • $.* - все поля корневого уровня

Производительность

  • Filter - применяйте рано для уменьшения объема данных
  • Select - уменьшает размер сообщений и повышает производительность
  • Flatten - может увеличить количество сообщений, используйте осторожно
  • Router - создает дополнительные подключения, минимизируйте количество маршрутов

SnakeCase

Преобразует все ключи JSON-объекта в формат snake_case. Полезно для нормализации имен полей при интеграции с системами, использующими snake_case (например, PostgreSQL, Python API).

Конфигурация

transformations:
  - type: snakeCase
    config:
      # Рекурсивно преобразовывать вложенные объекты (опционально, по умолчанию: false)
      deep: true

Примеры

Простое преобразование

transformations:
  - type: snakeCase
    config:
      deep: false

Входное сообщение:

{
  "firstName": "John",
  "lastName": "Doe",
  "userName": "johndoe",
  "isActive": true,
  "itemCount": 42
}

Выходное сообщение:

{
  "first_name": "John",
  "last_name": "Doe",
  "user_name": "johndoe",
  "is_active": true,
  "item_count": 42
}

Рекурсивное преобразование

transformations:
  - type: snakeCase
    config:
      deep: true

Входное сообщение:

{
  "firstName": "John",
  "address": {
    "streetName": "Main St",
    "houseNumber": 123,
    "zipCode": "12345"
  },
  "items": [
    {
      "itemName": "Product",
      "itemPrice": 99.99
    }
  ]
}

Выходное сообщение:

{
  "first_name": "John",
  "address": {
    "street_name": "Main St",
    "house_number": 123,
    "zip_code": "12345"
  },
  "items": [
    {
      "item_name": "Product",
      "item_price": 99.99
    }
  ]
}

Преобразование PascalCase

transformations:
  - type: snakeCase
    config:
      deep: false

Входное сообщение:

{
  "FirstName": "John",
  "LastName": "Doe",
  "UserID": 123
}

Выходное сообщение:

{
  "first_name": "John",
  "last_name": "Doe",
  "user_id": 123
}

Особенности

  • Преобразует camelCasesnake_case
  • Преобразует PascalCasesnake_case
  • Обрабатывает последовательные заглавные буквы (например, XMLHttpRequestxml_http_request)
  • Сохраняет уже существующие ключи в snake_case без изменений
  • При deep: false преобразует только ключи верхнего уровня
  • При deep: true рекурсивно преобразует все вложенные объекты и массивы

CamelCase

Преобразует все ключи JSON-объекта в формат CamelCase (PascalCase). Полезно для нормализации имен полей при интеграции с системами, использующими CamelCase (например, Java, C# API).

Конфигурация

transformations:
  - type: camelCase
    config:
      # Рекурсивно преобразовывать вложенные объекты (опционально, по умолчанию: false)
      deep: true

Примеры

Простое преобразование

transformations:
  - type: camelCase
    config:
      deep: false

Входное сообщение:

{
  "first_name": "John",
  "last_name": "Doe",
  "user_name": "johndoe",
  "is_active": true,
  "item_count": 42
}

Выходное сообщение:

{
  "FirstName": "John",
  "LastName": "Doe",
  "UserName": "johndoe",
  "IsActive": true,
  "ItemCount": 42
}

Рекурсивное преобразование

transformations:
  - type: camelCase
    config:
      deep: true

Входное сообщение:

{
  "first_name": "John",
  "address": {
    "street_name": "Main St",
    "house_number": 123,
    "zip_code": "12345"
  },
  "items": [
    {
      "item_name": "Product",
      "item_price": 99.99
    }
  ]
}

Выходное сообщение:

{
  "FirstName": "John",
  "Address": {
    "StreetName": "Main St",
    "HouseNumber": 123,
    "ZipCode": "12345"
  },
  "Items": [
    {
      "ItemName": "Product",
      "ItemPrice": 99.99
    }
  ]
}

Преобразование одиночных слов

transformations:
  - type: camelCase
    config:
      deep: false

Входное сообщение:

{
  "name": "John",
  "id": 123
}

Выходное сообщение:

{
  "Name": "John",
  "Id": 123
}

Особенности

  • Преобразует snake_caseCamelCase
  • Все слова начинаются с заглавной буквы (PascalCase)
  • Сохраняет уже существующие ключи в CamelCase без изменений
  • При deep: false преобразует только ключи верхнего уровня
  • При deep: true рекурсивно преобразует все вложенные объекты и массивы

DebeziumUnwrap

Преобразует событие Debezium из Kafka (payload.before/after/op) в «плоское» сообщение строки, удобное для обычных трансформаций и приёмников.

Упрощённый аналог SMT Debezium ExtractNewRecordState (NRSE) для JSON-envelope (без schema-aware Connect records).

Конфигурация

transformations:
  - type: debeziumUnwrap
    config:
      # Преобразовывать tombstone Kafka (пустой value) в operation=delete из metadata.key (опционально, по умолчанию: false)
      inferDeleteFromTombstone: true

      # Копировать payload.source.* в metadata как source_<field> (опционально, по умолчанию: false)
      includeSourceInMetadata: true

      # Операция для snapshot-событий Debezium (op="r"): insert (по умолчанию) или update
      snapshotOperation: insert

      # Записать NRSE-поля __op / __deleted в payload строки (опционально, по умолчанию: false)
      addOperationFields: true

      # Скопировать выбранные ключи payload.source в строку как source_<key> (опционально, по умолчанию: выкл.)
      addSourceFields: [table, lsn, ts_ms]

Поведение

  • op=c -> тело из payload.after, metadata.operation=insert
  • op=u -> тело из payload.after, metadata.operation=update
  • op=r -> тело из payload.after, metadata.operation=insert (или update, если snapshotOperation: update)
  • op=d -> тело из payload.before, metadata.operation=delete
  • Если в сообщении нет Debezium envelope (payload.op) — сообщение проходит без изменений.
  • metadata.operation остаётся insert / update / delete даже при addOperationFields (soft-delete sinks продолжают читать metadata).

При addOperationFields: true:

  • в строку добавляется __op = исходный Debezium op (c / u / d / r)
  • добавляется __deleted = "true" на delete (включая tombstone-inferred), иначе "false" (строки, удобно для JDBC)

При непустом addSourceFields:

  • для каждого ключа k, если есть payload.source[k], пишется source_<k> в строку
  • отсутствующие ключи пропускаются; независимо от includeSourceInMetadata

Паритет с Debezium ExtractNewRecordState

Возможность Debezium NRSE DataFlow debeziumUnwrap
Flatten envelope → row Да Да (after / before по op)
__op / __deleted в payload Да Opt-in через addOperationFields
Source-поля в payload add.source.fields Opt-in через addSourceFields
Source в headers/metadata Headers Opt-in через includeSourceInMetadata
Drop tombstones drop.tombstones По умолчанию drop; сохранить через inferDeleteFromTombstone (+ filter)
MongoDB unwrap Отдельный SMT Не поддерживается
Schema-aware Connect records Да Только JSON envelope
Soft-delete column rewrite Config flags Sink softDeleteColumn + metadata.operation

Tombstone Kafka

Если inferDeleteFromTombstone: true и value сообщения пустой, трансформация пытается разобрать metadata.key как JSON (обычно ключ Debezium) и сформировать delete-сообщение.

Если ключ отсутствует или невалиден, сообщение отбрасывается (без ошибки).

При addOperationFields: true inferred tombstone получает __op: "d" и __deleted: "true".

Совместимость с PostgreSQL sink

Для operation=delete текущий PostgreSQL sink использует soft delete (через softDeleteColumn). Для корректной обработки удалений включайте softDeleteColumn в sink.config.

Пример Debezium -> PostgreSQL

transformations:
  - type: debeziumUnwrap
    config:
      inferDeleteFromTombstone: true
      includeSourceInMetadata: true
      snapshotOperation: insert
      addOperationFields: true
      addSourceFields: [table, lsn]

ReplaceField

Переименовывает поля и опционально фильтрует payload через include / exclude. В отличие от select, include сохраняет вложенную структуру (без сплющивания ключей). Удобно при миграции с Kafka Connect ReplaceField.

Конфигурация

transformations:
  - type: replaceField
    config:
      # Переименования oldPath:newPath (опционально)
      renames:
        - oldName:newName
        - key.sku:sku

      # Оставить только эти пути, сохраняя вложенность (опционально; взаимоисключающе с exclude)
      include:
        - user.id
        - status

      # Удалить эти пути (опционально; взаимоисключающе с include)
      # exclude:
      #   - password

Нужен хотя бы один из renames, include или exclude. include и exclude нельзя задавать вместе.

Примеры

Переименование полей

transformations:
  - type: replaceField
    config:
      renames:
        - oldName:newName
        - key.sku:sku

Входное сообщение:

{
  "oldName": "value",
  "key": { "sku": "12345" },
  "other": "keep"
}

Выходное сообщение:

{
  "newName": "value",
  "sku": "12345",
  "other": "keep"
}

Include без сплющивания

transformations:
  - type: replaceField
    config:
      include:
        - user.id
        - user.name
        - status

Входное сообщение:

{
  "user": { "id": 1, "name": "John", "email": "john@example.com" },
  "status": "active",
  "extra": "drop-me"
}

Выходное сообщение:

{
  "user": { "id": 1, "name": "John" },
  "status": "active"
}

HeadersToPayload

Копирует заголовки Kafka (или другого source) из Metadata["headers"] в поля JSON payload. Удобно для проброса tracing ID, tenant-ключей и других header-значений в тело сообщения перед sinks, которые не сохраняют Kafka headers.

Kafka source заполняет Metadata["headers"] как map[string]string, если у записи есть headers.

Конфигурация

transformations:
  - type: headersToPayload
    config:
      # маппинги headerName:fieldPath (обязательно; минимум один)
      mappings:
        - X-Request-Id:requestId
        - X-Language:metadata.language

Отсутствующие headers пропускаются (payload для этого маппинга не меняется). Пути полей поддерживают вложенность и опциональный префикс $..

Примеры

Копирование headers в payload

transformations:
  - type: headersToPayload
    config:
      mappings:
        - X-Request-Id:requestId
        - X-Language:metadata.language

Входное сообщение:

{
  "data": "value"
}

Headers: X-Request-Id=req-123, X-Language=en

Выходное сообщение:

{
  "data": "value",
  "requestId": "req-123",
  "metadata": {
    "language": "en"
  }
}

StructFlatten

Сплющивает вложенные JSON-объекты в одноуровневую карту (1→1). В отличие от flatten (который разворачивает массив в N сообщений), structFlatten оставляет одно сообщение и склеивает пути ключей разделителем. Массивы сохраняются как значения и не индексируются. Удобно при миграции с Kafka Connect Flatten.

Корень не-объект (массив, примитив) и не-JSON payload проходят без изменений. Пустые вложенные объекты {} не порождают ключей. Вложенность глубже 64 уровней возвращает ошибку трансформации.

Конфигурация

transformations:
  - type: structFlatten
    config:
      # Разделитель сегментов вложенных ключей (опционально, по умолчанию ".")
      delimiter: "."

Пустой config: {} валиден и использует ".". Явно пустой delimiter: "" отклоняется.

Примеры

Разделитель «точка» (по умолчанию)

transformations:
  - type: structFlatten
    config:
      delimiter: "."

Входное сообщение:

{
  "content": {
    "id": 42,
    "name": {
      "first": "David",
      "middle": null,
      "last": "Wong"
    },
    "tags": ["a", "b"]
  },
  "active": true
}

Выходное сообщение:

{
  "content.id": 42,
  "content.name.first": "David",
  "content.name.middle": null,
  "content.name.last": "Wong",
  "content.tags": ["a", "b"],
  "active": true
}

Разделитель «подчёркивание» (имена для JDBC / Avro)

transformations:
  - type: structFlatten
    config:
      delimiter: "_"

Ключи на выходе: content_id, content_name_first, content_name_middle, content_name_last, content_tags, active.

Типичная CDC-цепочка: debeziumUnwrapstructFlattensnakeCase.

ExtractField

Заменяет payload сообщения значением одного поля (аналог Kafka Connect ExtractField$Value). Кардинальность всегда 1→1. Metadata сохраняется. Не-JSON и отсутствующий путь проходят без изменений. Новый корень может быть объектом, массивом, примитивом или JSON null.

Конфигурация

transformations:
  - type: extractField
    config:
      # JSONPath к полю, которое становится новым корнем (обязательно)
      field: payload.after   # или $.payload.after

Примеры

Распаковка вложенного payload

transformations:
  - type: extractField
    config:
      field: payload.after

Входное сообщение:

{"payload":{"after":{"id":1}}}

Выходное сообщение:

{"id":1}

Извлечение примитива или массива

transformations:
  - type: extractField
    config:
      field: items

{"items":[1,2,3]}[1,2,3]

Типичная цепочка перед flatten/cast: extractFieldstructFlattencast.

HoistField

Оборачивает весь JSON payload под одним ключом верхнего уровня (обратно к extractField). Удобно для нормализации обёрток до/после CDC. Не-JSON не меняется. Имя ключа — простое, без точек (не JSONPath).

Конфигурация

transformations:
  - type: hoistField
    config:
      # Ключ-обёртка верхнего уровня (обязательно; без точек)
      field: record

Примеры

Обернуть строку

transformations:
  - type: hoistField
    config:
      field: record

Входное сообщение:

{"id":1}

Выходное сообщение:

{"record":{"id":1}}

Round-trip с extractField

transformations:
  - type: hoistField
    config:
      field: record
  - type: extractField
    config:
      field: record

Восстанавливает исходный payload.

Cast

Приводит значения полей к заданным типам по JSONPath. Удобно после debeziumUnwrap / structFlatten перед JDBC и другими sinks со строгой схемой. Кардинальность 1→1. Metadata не меняется. Не-JSON проходит без изменений.

Отсутствующий путь пропускается (не ошибка). Неуспешный cast существующего значения — ошибка transform: processor логирует, инкрементирует dataflow_transformer_errors_total и пропускает сообщение (в sink не пишется). Оператор не перезапускается. В Kafka failed offset этим сообщением не ack'ается; последующая успешная запись может продвинуть commit дальше (skip-on-error, не бесконечный retry на том же offset).

Конфигурация

transformations:
  - type: cast
    config:
      # Карта JSONPath → целевой тип (обязательна, непустая)
      # Типы: string | int64 | float64 | bool | null
      spec:
        id: int64
        amount: float64
        active: bool
        note: string
        deleted_at: null    # принудительно JSON null

Правила конверсии

Цель Принимает Ошибка на
string скаляры (числа, bool, строки) object, array, JSON null
int64 целые, целые float, числовые строки дробные числа, нечисловые значения
float64 числа, числовые строки нечисловые значения
bool bool; строки true/false (без учёта регистра); числа 0/1 прочие значения
null любое существующее значение → JSON null

Примеры

После unwrap / flatten

transformations:
  - type: debeziumUnwrap
  - type: cast
    config:
      spec:
        id: int64
        amount: float64
        active: bool

Входное сообщение:

{"id":"1","amount":"9.99","active":"true"}

Выходное сообщение:

{"id":1,"amount":9.99,"active":true}

Типичная цепочка: debeziumUnwraptimezonecast или extractFieldstructFlattencast.

Timezone

Конвертирует перечисленные temporal-поля в целевую IANA-зону или фиксированный UTC offset (±HH:MM). Полезно после debeziumUnwrap перед JDBC sinks. Кардинальность 1→1. Metadata не меняется. Не-JSON проходит без изменений.

Не связана с spec.maintenance.timezone и отличается от transform timestamp (вставка текущего wall-clock времени).

Отсутствующие поля и JSON null пропускаются. Непарсабельные значения — ошибка transform: processor логирует, инкрементирует dataflow_transformer_errors_total и пропускает сообщение (те же skip-on-error семантики, что у cast).

Конфигурация

transformations:
  - type: timezone
    config:
      # Целевая IANA TZ или UTC offset (обязательно)
      timezone: Europe/Moscow
      # Поля для конвертации (обязательно, непустой список JSONPath)
      fields: [created_at, updated_at]
      # Исходная TZ, если у значения нет offset (опционально, по умолчанию: UTC)
      sourceTimezone: UTC
      # Формат вывода (опционально, по умолчанию: RFC3339Nano). Также: RFC3339, UnixMilli
      format: RFC3339

Формы входа

Вход Поведение
Строка RFC3339 / RFC3339Nano (с offset или без) Парсится; без offset — в sourceTimezone
Epoch (число или numeric string) Миллисекунды если \|n\| >= 1e12, иначе секунды
Отсутствующий путь / JSON null Skip
Прочее Ошибка transform

Примеры

После debeziumUnwrap

transformations:
  - type: debeziumUnwrap
  - type: timezone
    config:
      timezone: Europe/Moscow
      fields: [created_at, updated_at]
      format: RFC3339
  - type: cast
    config:
      spec:
        id: int64

Входное сообщение:

{"id":1,"created_at":"2024-01-15T12:00:00Z"}

Выход после timezone (Europe/Moscow):

{"id":1,"created_at":"2024-01-15T15:00:00+03:00"}

InsertField

Вставляет или перезаписывает JSON-поля литералами, плейсхолдерами metadata записи, текущим временем или сырым JSON. Аналог Kafka Connect InsertField$Value с metadata DataFlow и префиксом json:. Кардинальность 1→1. Metadata не меняется. Не-JSON проходит без изменений.

Не путать с timestamp (одно поле now) и headersToPayload (только headers → body).

Конфигурация

transformations:
  - type: insertField
    config:
      fields:
        pipeline: "orders-cdc"                 # литерал
        source_topic: "${metadata.topic}"      # плейсхолдер metadata
        source_partition: "${metadata.partition}"
        source_offset: "${metadata.offset}"
        source_timestamp: "${metadata.timestamp}"
        ingested_at: "${now}"                  # wall-clock RFC3339
        "flags.reprocessed": "json:false"      # сырое JSON-значение

Синтаксис значений

Значение Результат
любая другая строка JSON-строка
${metadata.<key>} строковое представление Metadata[key]; нет/nil → ""
${now} time.Now() в формате RFC3339
json:<raw> распарсенный JSON как нативное значение (false, 42, объекты, массивы)

Типичные ключи metadata у Kafka source: topic, partition, offset, timestamp. Невалидный json: — ошибка transform (сообщение пропускается).

Примеры

Обогащение контекстом пайплайна

transformations:
  - type: insertField
    config:
      fields:
        pipeline: "orders-cdc"
        source_topic: "${metadata.topic}"
        ingested_at: "${now}"
        "flags.reprocessed": "json:false"
  - type: snakeCase
    config:
      deep: true

Вход: {"id":1,"Name":"A"} с metadata topic=raw.events
После insertField: добавлены pipeline, source_topic, ingested_at, flags.reprocessed
После snakeCase: ключи в snake_case

Ограничения

  • Filter / Router / when: общий DSL условий — истинность, == / !=, metadata.*, && / || / !, скобки. Отсутствующее поле не проходит условие. Groovy/JS/CEL не поддерживаются. when пропускает шаг (passthrough); filter отбрасывает сообщение.
  • Flatten: работает только с массивами (в т.ч. в обёртке Avro с ключом array), не с произвольными объектами. Для сплющивания вложенных объектов используйте structFlatten.
  • StructFlatten: сплющивание объектов 1→1; массивы сохраняются как значения (сначала array-flatten, если нужен explode). Максимальная глубина вложенности — 64.
  • ExtractField: отсутствующий путь и не-JSON → passthrough; корень может стать не-объектом (следующие object-only transforms сделают passthrough).
  • HoistField: ключ-обёртка — простое имя верхнего уровня (без точек); оборачивает любой JSON, включая массивы и примитивы.
  • Cast: отсутствующий путь → skip; неуспешный cast существующего значения → ошибка transform (сообщение пропускается, в sink не пишется; Kafka offset может уйти вперёд при следующем успехе). Типы: string, int64, float64, bool, null.
  • Timezone: только перечисленные fields; отсутствующее/null → skip; непарсабельное → ошибка transform (те же skip-on-error, что у cast). Форматы: RFC3339, RFC3339Nano (по умолчанию), UnixMilli. Не связана с maintenance.timezone и transform timestamp.
  • InsertField: не-JSON → passthrough; нет metadata → пустая строка; невалидный json: → ошибка transform. Отличается от timestamp и headersToPayload.
  • Select: результат всегда плоский; ключом становится последний сегмент JSONPath.
  • ReplaceField: include сохраняет вложенность (в отличие от select); include и exclude взаимоисключающие.
  • HeadersToPayload: нужны headers в Metadata["headers"] (Kafka source их выставляет); отсутствующие headers пропускаются; не-JSON payload не меняется.
  • SnakeCase и CamelCase: работают только с валидным JSON; бинарные данные возвращаются без изменений.
  • DebeziumUnwrap: поддерживает Debezium JSON envelope (payload.op/before/after). Для operation=delete в PostgreSQL sink требуется softDeleteColumn, иначе событие может быть обработано как обычная вставка/апдейт.