В этом разделе рассматривается пример создания инстанса Kafka (ciphertext access и SASL_SSL) и доступа к нему на клиенте (private network, внутри виртуальной частной облачной сети (VPC)) для производства и потребления сообщений, чтобы быстро начать работу с Distributed Message Service (DMS) for Kafka.
Figure 1 Процедура использования DMS for Kafka

Инстанс Kafka работает в виртуальной частной облачной сети (VPC). Перед созданием инстанса Kafka убедитесь, что VPC доступна.
После создания инстанса Kafka загрузите и установите открытый клиент Kafka на ваш ECS перед производством и потреблением сообщений.
При создании инстанса Kafka вы можете выбрать спецификацию и количество, а также включить ciphertext access и SASL_SSL.
При подключении к инстансу Kafka с включённым SASL_SSL используется SASL для аутентификации. Данные шифруются SSL‑сертификатами для передачи с высоким уровнем безопасности.
Топики хранят сообщения, созданные продюсерами и получаемые потребителями.
В этом разделе используется пример создания топика в консоли.
Перед подключением к инстансу Kafka с включённым SASL_SSL загрузите сертификат и настройте соединение в файле конфигурации клиента.
Для обеспечения детального управления вашими облачными ресурсами создайте группы пользователей и пользователей Identity and Access Management (IAM) и предоставьте указанным пользователям соответствующие разрешения. Для получения дополнительной информации см. Creating an IAM User and Granting DMS for Kafka Permissions.
Перед созданием экземпляра Kafka убедитесь, что VPC и подсеть доступны. Подробную информацию о том, как создать VPC и подсеть, см. Creating a VPC.
VPC должна быть создана в том же регионе, где находится экземпляр Kafka.
Перед созданием экземпляра Kafka убедитесь, что группа безопасности доступна. Подробную информацию о том, как создать группу безопасности, см. Creating a Security Group.
Группа безопасности должна быть создана в том же регионе, где находится экземпляр Kafka.
Чтобы подключиться к экземплярам Kafka, добавьте правила группы безопасности, описанные в Table 1. Другие правила могут быть добавлены в соответствии с требованиями площадки.
Направление | Протокол | Порт | Исходный адрес | Описание |
|---|---|---|---|---|
Inbound | TCP | 9093 | 0.0.0.0/0 | Доступ к экземпляру Kafka через частную сеть внутри VPC (в зашифрованном виде) |
После создания группы безопасности у неё есть правило входящего трафика по умолчанию, позволяющее взаимодействовать ECS внутри группы безопасности, и правило исходящего трафика по умолчанию, позволяющее весь исходящий трафик. Если вы получаете доступ к вашему экземпляру Kafka через частную сеть внутри VPC, добавлять правила, описанные в Table 1, не требуется.
В этом разделе в качестве клиента используется Linux Elastic Cloud Server (ECS). Перед созданием экземпляра Kafka создайте ECS с EIP, установите JDK, настройте переменные окружения и загрузите клиент Kafka с открытым исходным кодом. Подробную информацию о создании ECS см. в Creating an ECS. Если у вас уже есть доступный ECS, пропустите этот шаг. Используйте Oracle JDK вместо JDK по умолчанию в ECS (например, OpenJDK), поскольку JDK по умолчанию в ECS может быть неподходящим. Получите Oracle JDK версии 1.8.111 или новее с Oracle's official website. Замените jdk-8u321-linux-x64.tar.gz на вашу версию JDK. Измените /root/jdk1.8.0_321 на путь, где вы установили JDK. Если возвращается следующее сообщение, JDK установлен.
в левом верхнем углу, нажмите Elastic Cloud Server в разделе Computing, а затем создайте ECS.
Параметр | Описание |
|---|---|
Режим биллинга | Выберите Pay-per-use, который является постоплатным режимом. Вы можете платить после использования сервиса, и вам будет выставлен счет за продолжительность использования. Плата рассчитывается в секундах и выставляется почасово. |
Region | DMS for Kafka в разных регионах не могут взаимодействовать друг с другом через интранет. Выберите ближайшее расположение для низкой задержки и быстрого доступа. Выберите RU-Moscow. |
Project | Проекты изолируют вычислительные, хранилищные и сетевые ресурсы в разных географических регионах. Для каждого региона доступен предустановленный проект. Выберите RU-Moscow (по умолчанию). |
AZ | AZ — это физический регион, где ресурсы используют независимое электропитание и сети. AZ физически изолированы, но соединены через внутреннюю сеть. Выберите AZ1, AZ2 и AZ3. |
Instance Name | Вы можете задать имя, соответствующее правилам: 4–64 символа; начинается с буквы; может содержать только буквы, цифры, дефисы (-) и подчёркивания (_). Введите kafka-test. |
Enterprise Project | Этот параметр предназначен для корпоративных пользователей. Enterprise project управляет ресурсами проекта в группах. Enterprise projects логически изолированы. Выберите default. |
Version | Версия Kafka. Не может быть изменена после создания инстанса. Выберите 2.7. |
Broker Flavor | Выберите broker flavor по требованию. Выберите kafka.2u4g.cluster. |
Brokers | Укажите количество брокеров по требованию. Введите 3. |
Storage Space per Broker | Выберите тип диска и укажите размер диска по требованию. Примечание: Вместимость диска может быть установлена только целым кратным 100. Общий объём хранения = Storage space per broker × Broker quantity. После создания инстанса изменить тип диска нельзя. Выберите Ultra-high I/O и введите 100. |
Capacity Threshold Policy | Выберите Automatically delete: Когда диск достигает порога заполнения (95 %), сообщения могут продолжать создаваться и потребляться, но первые 10 % сообщений будут удалены для обеспечения достаточного свободного места. Используйте эту политику для сервисов, не допускающих прерываний. Однако данные могут быть потеряны. |
Параметр | Подпараметр | Описание |
|---|---|---|
Private Network Access | Plaintext Access | Skip it. |
Ciphertext Access | Когда этот параметр включён, требуется аутентификация SASL при подключении клиента к экземпляру Kafka.
| |
Private IP Addresses | Выберите Auto: Система автоматически назначает IP-адреса из подсети. | |
Public Network Access | - | Пропустить. |
Создание instance занимает от 3 до 15 минут. В течение этого периода статус instance — Creating.
Instances, которые не удалось создать, не занимают другие ресурсы.
Figure 2 Адреса Kafka instance (частная сеть) для доступа внутри VPC

Параметр | Описание |
|---|---|
Имя темы | Укажите имя, содержащее от 3 до 200 символов, начинающееся с буквы или символа подчёркивания (_), и содержащее только буквы, цифры, точки (.), дефисы (-) и подчёркивания (_). Имя должно отличаться от предустановленных тем:
Cannot be changed once the topic is created. Введите topic-01. |
Разделы | Если количество разделов совпадает с количеством потребителей, чем больше разделов, тем выше параллельность потребления. Введите 3. |
Реплики | Данные автоматически сохраняются в каждой реплике. Если один брокер Kafka выходит из строя, данные остаются доступными. Большее количество реплик обеспечивает более высокую надёжность. Введите 3. |
Время устаревания (ч) | Как долго сообщения будут сохраняться в теме. Сообщения старше этого периода не могут быть потреблены. Они будут удалены, и их более нельзя будет потреблять. Введите 72. |
Synchronous Replication | Пропустить. Когда эта опция отключена, лидирующие реплики независимы от синхронизации реплик‑подписчиков. Они получают сообщения и записывают их в локальные журналы, затем сразу отправляют успешно записанные клиенту. |
Synchronous Flushing | Пропустить. Когда эта опция отключена, сообщения создаются и хранятся в памяти вместо немедленной записи на диск. |
Message Timestamp | Выберите CreateTime: время, когда продюсер создал сообщение. |
Max. Message Size (bytes) | Максимальный размер пакетной обработки, разрешённый Kafka. Если сжатие сообщений включено в файле конфигурации клиента или в коде продюсеров, этот параметр указывает размер после сжатия. Введите 10,485,760. |
Description | Пропустить. |
Чтобы получить сертификат: В консоли Kafka нажмите на экземпляр Kafka, чтобы перейти на страницу Basic Information. Нажмите Download рядом с SSL Certificate в области Connection. Распакуйте пакет, чтобы получить файл сертификата client.jks.
/root — путь для хранения сертификата. При необходимости измените его на фактический путь.
cd kafka_2.12-2.7.2/config
sasl.jaas.config=org.apache.kafka.common.security.plain.PlainLoginModule required \username="**********" \password="**********";sasl.mechanism=PLAINsecurity.protocol=SASL_SSLssl.truststore.location={ssl_truststore_path}ssl.truststore.password=dms@kafkassl.endpoint.identification.algorithm=
Описание:
cd ../bin
./kafka-console-producer.sh --broker-list {connection address} --topic {topic name} --producer.config ../config/producer.properties
Описание:
Например, 192.xxx.xxx.xxx:9093, 192.xxx.xxx.xxx:9093, 192.xxx.xxx.xxx:9093 являются адресами подключения Kafka instance.
После выполнения этой команды вы можете отправлять сообщения в Kafka instance, вводя запрашиваемую информацию и нажимая Enter. Каждая строка содержимого будет отправлена как сообщение.
[root@ecs-kafka bin]#./kafka-console-producer.sh --broker-list 192.xxx.xxx.xxx:9093,192.xxx.xxx.xxx:9093,192.xxx.xxx.xxx:9093 --topic topic-01 --producer.config ../config/producer.properties>Hello>DMS>Kafka!>^C[root@ecs-kafka bin]#
Нажмите Ctrl+C, чтобы отменить.
./kafka-console-consumer.sh --bootstrap-server {connection address} --topic {topic name} --from-beginning --consumer.config ../config/consumer.properties
Описание:
Пример:
[root@ecs-kafka bin]# ./kafka-console-consumer.sh --bootstrap-server 192.xxx.xxx.xxx:9093,192.xxx.xxx.xxx:9093,192.xxx.xxx.xxx:9093 --topic topic-01 --from-beginning --consumer.config ../config/consumer.propertiesHelloKafka!DMS^CProcessed a total of 3 messages[root@ecs-kafka bin]#
Нажмите Ctrl+C, чтобы отменить.