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сообщения (например Kafkatopic,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
Рекомендации по порядку
- Flatten должен быть первым, если нужно развернуть массивы
- Filter применяйте рано, чтобы уменьшить объем обрабатываемых данных
- SnakeCase/CamelCase применяйте после Select/Remove, но перед отправкой в приемник
- Mask/Remove применяйте перед Select для безопасности
- Select применяйте в конце для финальной очистки
- Timestamp можно применять в любом месте, но обычно в начале или конце
- 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
}
Особенности
- Преобразует
camelCase→snake_case - Преобразует
PascalCase→snake_case - Обрабатывает последовательные заглавные буквы (например,
XMLHttpRequest→xml_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_case→CamelCase - Все слова начинаются с заглавной буквы (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=insertop=u-> тело изpayload.after,metadata.operation=updateop=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-цепочка: debeziumUnwrap → structFlatten → snakeCase.
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: extractField → structFlatten → cast.
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}
Типичная цепочка: debeziumUnwrap → timezone → cast или extractField → structFlatten → cast.
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и transformtimestamp. - 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, иначе событие может быть обработано как обычная вставка/апдейт.