Объединение Apache Kafka с Greenplum
Greenplum Stream Server (GPSS) - это инструмент ETL (извлечение, преобразование, загрузка). Экземпляр сервера GPSS получает потоковые данные от одного или нескольких клиентов, используя внешние таблицы Greenplum, доступные для чтения, преобразования и вставки данных в целевую таблицу Greenplum. Источник данных и формат данных зависят от каждого конкретного клиента.
Сервер потоков Greenplum включает утилиту командной строки GPSS. Когда Вы запускаете его, Вы запускаете экземпляр GPSS, который ожидает клиентских данных неопределенное время.
Сервер Greenplum Stream Server также включает утилиту командной строки gpsscli, клиентский инструмент для отправки для отправки заданий на загрузку данных в экземпляр GPSS и управления этими заданиями.
Ограничения
Greenplum Stream Server не поддерживает загрузку данных из нескольких тем Kafka в одну и ту же таблицу Greenplum. Все задания будут зависать, если GPSS столкнется с такой ситуацией.
Приступим к работе.
Немного подготовки
Со стороны GPMaster
- Подготовьте базу данных + таблицу, в которую будут поступать данные из kafka
testdb=# CREATE TABLE json_from_kafka( customer_id int8, month int4, amount_paid decimal(9,2) );
- Зарегистрируйте расширение GPSS
Вы должны зарегистрировать расширение Greenplum Stream Server в каждой базе данных, в которой Вы будете использовать его для записи данных в таблицы Greenplum.
$ ssh gpadmin@gpmaster gpmaster$ . /usr/local/greenplum-db/greenplum_path.sh mdw$ psql -d testdb testdb=# CREATE EXTENSION gpss;
Для каждой базы данных, в которую Greenplum Stream Server будет записывать клиентские данные, выполните шаги 3 и 4.
- Настройка сервера Greenplum Stream Server через наш мастер
- Наш сервер потоков будет объединять в рамках одного процесса и gpss listener, и gpfdist.
Создайте файл «gpss_config.json», который будет отвечать за настройку нашего сервиса GPSS
Адрес хоста по умолчанию— localhost
Пример:
{
“ListenAddress”: {
“Host”: “”,
“Port”: 50007,
“SSL”: false
},
“Gpfdist”: {
“Host”: “”,
“Port”: 8319
}
}
gpmaster$ gpss gpss_config.json — log-dir . &
- Теперь у нас есть оба слушателя 'gpfdist' для отправки данных на наши сегментные узлы + слушатель 'gpss' для получения данных с клиентского узла Kafka.
- Создание файла ‘jsonload_cfg.yaml’
В этом файле мы настраиваем наш источник данных и указываем, где его необходимо разместить, в какой базе данных и в какой таблице.
https://gpdb.docs.pivotal.io/5110/greenplum-kafka/load-json-example.html
Пример:
Со стороны Клиента
1. Установите Kafka
В этом руководстве мы будем использовать Kafka в качестве брокера сообщений или ETL-сервера. Он будет получать топики или данные и отправлять их нашему gpss-слушателю на главном узле.
2. Убедитесь в том, что у Вас есть маршрутизация между сервером kafka и сервером gpmaster на порту 9092
3. Создайте топик Kafka
Greenplum имеет ограничение 1:1 между топиком и таблицей. Он не может подключить несколько топиков к одной таблице.
https://gpdb.docs.pivotal.io/5110/greenplum-kafka/load-json-example.html
Пример:
kafkahost$ $KAFKA_INSTALL_DIR/bin/kafka-topics.sh — create \ — zookeeper localhost:2181 — replication-factor 1 — partitions 1 \ — topic topic_json_gpkafka
4. Добавьте данный в новый топик Kafka
kafkahost$ vi sample_data.json
Скопируйте и вставьте:
{ “cust_id”: 1313131, “month”: 12, “expenses”: 1313.13 }
{ “cust_id”: 3535353, “month”: 11, “expenses”: 761.35 }
{ “cust_id”: 7979797, “month”: 10, “expenses”: 4489.00 }
{ “cust_id”: 7979797, “month”: 11, “expenses”: 18.72 }
{ “cust_id”: 3535353, “month”: 10, “expenses”: 6001.94 }
{ “cust_id”: 7979797, “month”: 12, “expenses”: 173.18 }
{ “cust_id”: 1313131, “month”: 10, “expenses”: 492.83 }
{ “cust_id”: 3535353, “month”: 12, “expenses”: 81.12 }
{ “cust_id”: 1313131, “month”: 11, “expenses”: 368.27 }
5. Отправьте данные в созданную нами тему Kafka
kafkahost$ $KAFKA_INSTALL_DIR/bin/kafka-console-producer.sh \ — broker-list localhost:9092 \ — topic topic_json_gpkafka < sample_data.json
6. Убедитесь в том, что данные были добавлены
kafkahost$ $KAFKA_INSTALL_DIR/bin/kafka-console-consumer.sh \ — bootstrap-server localhost:9092 — topic topic_json_gpkafka \ — from-beginning
Отправка запроса
Со стороны мастера
gpmaster$ gpsscli submit — name kafkajson2gp — gpss-port 50007 ./jsonload_cfg.yaml
1. Пролистайте все задания
gpmaster$ gpsscli list — all — gpss-port 50007
2.
Начните выполнение задания
gpmaster$ gpsscli start kafkajson2gp — gpss-port 50007
3. Для того, чтобы остановить прием данных и вставить новые строки, остановите задание
gpmaster$ gpsscli stop kafkajson2gp — gpss-port 50007
4. Обратите внимание на вывод команды 'gpss'. Он должен выглядеть следующим образом:
… -[INFO]:- … Inserted 9 rows … -[INFO]:- … Rejected 0 rows
5. Просмотрите новый контент в таблице, которую мы только что создали
gpmaster$ psql -d testdb testdb=# SELECT * FROM json_from_kafka WHERE customer_id=’1313131' ORDER BY amount_paid; customer_id | month | amount_paid — — — — — — -+ — — — -+ — — — — — — - 1313131 | 11 | 368.27 1313131 | 10 | 492.83 1313131 | 12 | 1313.13






