Восстановление топика из S3#
Предупреждение
Поддерживается только SELECT * запрос
Класс коннектора#
io.lenses.streamreactor.connect.aws.s3.source.S3SourceConnector
Поддержка KCQL#
Для настройки выбираемых данных для бекапирования можно использовать SQL-подобный синтаксис — KCQL. Формат KCQL-синтаксиса для source-коннектора:
INSERT INTO <имя топика>
SELECT *
FROM <имя бакета>:<pathPrefix>
[BATCH=batch]
[STOREAS <формат данных>]
[LIMIT limit]
[PROPERTIES(
'свойство 1'=x,
'свойство 2'=x,
)]
В KCQL поддерживается экранирование для операторов INSERT INTO, SELECT * FROM и PARTITIONBY. Например, экранирование имени my-topic-with-hyphen, в котором используется дефис:
INSERT INTO `my-topic-with-hyphen`
SELECT *
FROM bucketAddress:pathPrefix
Для указания нескольких KCQL запросов в рамках одной конфигурации разделите запросы символом ;. Ограничение: для восстановления разных топиков из одного источника (пара <имя бакета>:<pathPrefix>) настройте отдельный экземпляр коннектора.
Настройка источника#
S3-источник задается в команде FROM. Коннектор считывает все объекты из указанного местоположения, с учетом параметров разделения и упорядочивания данных. Для каждого раздела данных должен быть задан свой коннектор.
Формат команды FROM:
FROM <имя бакета>:<pathPrefix>
//my-bucket-called-pears:my-folder-called-apples
Если:
данные в AWS не были записаны с помощью коннектора, настроенного на обход иерархии в бакете и загрузке на основании даты последней модификации объекта в корзине (
connect.s3.source.partition.extractor.regex=none);используется
LastModifiedсортировка (connect.s3.source.ordering.type=LastModified).
убедитесь, что объекты не поступают с опозданием, или используйте этап постобработки для их обработки.
Для загрузки в буквенно-цифровом порядке использйте AlphaNumeric порядок.
Настройка получателя#
Целевой Kafka топик указывается в команде INSERT INTO, в него запишутся все данные:
INSERT INTO my-apples-topic SELECT * FROM my-bucket-called-pears:my-folder-called-apples
Форматы объектов S3#
Форматы хранения, поддерживаемые коннектором::
JSON— коннектор считывает сообщения в формате JSON, где каждая строка представляет собой отдельную запись.Avro— коннектор считывает сообщения в формате Avro из S3 и преобразовывать их в собственный формат Kafka.Parquet— коннектор считывает сообщения в формате Parquet из S3 и преобразовывать их в собственный формат Kafka.Text— коннектор считывает строки текста, каждая из которых представляет собой отдельную запись.CSV— коннектор считывает строки текста, каждая из которых представляет собой отдельную запись.CSV_WithHeaders— коннектор считывает строки текста, каждая из которых представляет собой отдельную запись, пропуская строку заголовка.Bytes— коннектор считывает байты, и преобразовывать каждый объект в сообщение Kafka.
Используйте команду STOREAS для настройки формата хранения, примеры:
STOREAS `JSON`
STOREAS `Avro`
STOREAS `Parquet`
STOREAS `Text`
STOREAS `CSV`
STOREAS `CSV_WithHeaders`
STOREAS `Bytes`
Обработка текста#
При использовании текстового хранилища (тип объекта Text) коннектор предоставляет дополнительные параметры конфигурации для обработки контента.
Регулярное выражение#
В режиме регулярного выражения коннектор применяет шаблон регулярного выражения, и только если строка соответствует шаблону, она считается записью. Например, чтобы включить только строки, начинающиеся с цифры, можно использовать следующую конфигурацию:
connect.s3.kcql=insert into $kafka-topic select * from lensesio:regex STOREAS `text` PROPERTIES('read.text.mode'='regex', 'read.text.regex'='^[1-9].*')
Начальная и конечная строки#
В данном режиме коннектор считывает текстовое содержимое между указанными строками, включая их. Например, чтобы прочитать записи, в которых первая строка — «SSM», а последняя строка — пустая (’’), используйте конфигурацию:
connect.s3.kcql=insert into $kafka-topic select * from lensesio:multi_line STOREAS `text` PROPERTIES('read.text.mode'='startEndLine', 'read.text.start.line'='SSM', 'read.text.end.line'='')
Чтобы обрезать начальную и конечную строки, установите для свойства read.text.trim значение true:
connect.s3.kcql=insert into $kafka-topic select * from lensesio:multi_line STOREAS `text` PROPERTIES('read.text.mode'='startEndLine', 'read.text.start.line'='SSM', 'read.text.end.line'='', 'read.text.trim'='true')
Начальный и конечный тег#
В данном режиме коннектор считывает текстовое содержимое между указанными тегами, включая их. Этот режим полезен, когда одна строка текста в S3 соответствует нескольким сообщениям в Kafka. Например, чтобы считывать XML-записи, заключенные между тегами <SSM> и </SSM>, используйте конфигурацию:
connect.s3.kcql=insert into $kafka-topic select * from lensesio:xml STOREAS `text` PROPERTIES('read.text.mode'='startEndTag', 'read.text.start.tag'='<SSM>', 'read.text.end.tag'='</SSM>')
Последовательность#
Данный коннектор использует нулевое заполнение в именах объектов для обеспечения точного упорядочивания, используя оптимизации, предлагаемые S3 API, и гарантируя точную последовательность объектов.
При использовании источника S3 вместе с приемником S3 коннектор может использовать метод упорядочивания, обеспечивая правильную хронологическую обработку данных. Однако, в некоторых случаях, данные S3 генерируются приложениями, которые не поддерживают лексический порядок имен ключей объектов.
В таких случаях установите connect.s3.source.ordering.type = LastModified. Это гарантирует, что источник упорядочивает и обрабатывает данные по дате изменения объектов.
При использовании LastModified сортировки убедитесь, что объекты не поступают с опозданием, или используйте этап постобработки для их обработки.
Тротлинг#
Чтобы ограничить количество считываемых за один опрос ключей используйте KCQL-команду BATCH (по умолчанию — 1000):
BATCH = 100
Чтобы ограничить количество возвращаемых за одну операцию опроса строк используйте команду KCQL-LIMIT (по умолчанию — 10000):
LIMIT 10000
Пример конфигурации#
name=local-s3-source
connector.class=io.lenses.streamreactor.connect.aws.s3.source.S3SourceConnector
topic=test
connect.s3.custom.endpoint={{ ENDPOINT }}
connect.s3.vhost.bucket=true
connect.s3.aws.auth.mode=Credentials
connect.s3.aws.region=eu-west-2
connect.s3.aws.access.key=none
connect.s3.aws.secret.key=none
connect.s3.kcql=insert into test select * from test STOREAS `TEXT`