Восстановление топика из 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`