Необходимо отфильтровать записи Kafka на основе определенного ключевого слова - PullRequest
0 голосов
/ 13 марта 2019

У меня есть тема Кафки, в которой около 3 миллионов записей. Я хочу выбрать одну запись из этого, которая имеет определенный параметр. Я пытался сделать запрос, используя линзы, но не смог сформировать правильный запрос. ниже приведены записи 1 сообщения.

{
  "header": {
    "schemaVersionNo": "1",
  },
  "payload": {
    "modifiedDate": 1552334325212,
    "createdDate": 1552334325212,
    "createdBy": "A",
    "successful": true,
    "source_order_id": "3411976933214",
  }
}

Теперь я хочу отфильтровать запись с определенным source_order_id, но не могу найти правильный способ сделать это. Мы также пробовали использовать линзы Kafka Tool.

Пример запроса, который мы пробовали в линзах, приведен ниже:

SELECT * FROM `TEST`
WHERE _vtype='JSON' AND _ktype='BYTES'
AND _sample=2 AND _sampleWindow=200 AND payload.createdBy='fms'

Этот запрос работает, однако, если мы попробуем с идентификатором источника, как показано ниже, мы получим ошибку:

SELECT * FROM `TEST`
WHERE _vtype='JSON' AND _ktype='BYTES'
AND _sample=2 AND _sampleWindow=200 AND payload.source_order_id='3411976911924'



 Error : "Invalid syntax at line=3 and column=41.Invalid syntax for 'payload.source_order_id'. Field 'payload' resolves to primitive type STRING.

Потребление всех 3 миллионов записей через обычного потребителя, а затем повторение по нему не кажется мне оптимизированным подходом, поэтому поиск любых доступных решений для такого варианта использования.

1 Ответ

7 голосов
/ 13 марта 2019

Поскольку вы сказали, что открыты для других решений, вот одно, построенное с использованием KSQL .

Сначала давайте рассмотрим несколько примеров записей в теме источника:

$ kafkacat -P -b localhost:9092 -t TEST <<EOF
{ "header": { "schemaVersionNo": "1" }, "payload": { "modifiedDate": 1552334325212, "createdDate": 1552334325212, "createdBy": "A", "successful": true, "source_order_id": "3411976933214" } }
{ "header": { "schemaVersionNo": "1" }, "payload": { "modifiedDate": 1552334325412, "createdDate": 1552334325412, "createdBy": "B", "successful": true, "source_order_id": "3411976933215" } }
{ "header": { "schemaVersionNo": "1" }, "payload": { "modifiedDate": 1552334325612, "createdDate": 1552334325612, "createdBy": "C", "successful": true, "source_order_id": "3411976933216" } }
EOF

Используя KSQL, мы можем проверить тему с помощью PRINT:

ksql> PRINT 'TEST' FROM BEGINNING;
Format:JSON
{"ROWTIME":1552476232988,"ROWKEY":"null","header":{"schemaVersionNo":"1"},"payload":{"modifiedDate":1552334325212,"createdDate":1552334325212,"createdBy":"A","successful":true,"source_order_id":"3411976933214"}}
{"ROWTIME":1552476232988,"ROWKEY":"null","header":{"schemaVersionNo":"1"},"payload":{"modifiedDate":1552334325412,"createdDate":1552334325412,"createdBy":"B","successful":true,"source_order_id":"3411976933215"}}
{"ROWTIME":1552476232988,"ROWKEY":"null","header":{"schemaVersionNo":"1"},"payload":{"modifiedDate":1552334325612,"createdDate":1552334325612,"createdBy":"C","successful":true,"source_order_id":"3411976933216"}}

Затем объявимсхема по теме, которая позволяет нам запускать SQL против нее:

ksql> CREATE STREAM TEST (header STRUCT<schemaVersionNo VARCHAR>, 
                          payload STRUCT<modifiedDate BIGINT, 
                                        createdDate BIGINT, 
                                        createdBy VARCHAR, 
                                        successful BOOLEAN, 
                                        source_order_id VARCHAR>) 
                          WITH (KAFKA_TOPIC='TEST', 
                                VALUE_FORMAT='JSON');

Message
----------------
Stream created
----------------

Скажите KSQL работать со всеми данными в теме:

ksql> SET 'auto.offset.reset' = 'earliest';
Successfully changed local property 'auto.offset.reset' to 'earliest'. Use the UNSET command to revert your change.

И теперь мы можем выбратьвсе данные:

ksql> SELECT * FROM TEST;
1552475910106 | null | {SCHEMAVERSIONNO=1} | {MODIFIEDDATE=1552334325212, CREATEDDATE=1552334325212, CREATEDBY=A, SUCCESSFUL=true, SOURCE_ORDER_ID=3411976933214}
1552475910106 | null | {SCHEMAVERSIONNO=1} | {MODIFIEDDATE=1552334325412, CREATEDDATE=1552334325412, CREATEDBY=B, SUCCESSFUL=true, SOURCE_ORDER_ID=3411976933215}
1552475910106 | null | {SCHEMAVERSIONNO=1} | {MODIFIEDDATE=1552334325612, CREATEDDATE=1552334325612, CREATEDBY=C, SUCCESSFUL=true, SOURCE_ORDER_ID=3411976933216}
^CQuery terminated

или мы можем выборочно запросить их, используя запись -> для доступа к вложенным полям в схеме:

ksql> SELECT * FROM TEST 
        WHERE PAYLOAD->CREATEDBY='A';
1552475910106 | null | {SCHEMAVERSIONNO=1} | {MODIFIEDDATE=1552334325212, CREATEDDATE=1552334325212, CREATEDBY=A, SUCCESSFUL=true, SOURCE_ORDER_ID=3411976933214}

, а также выбрав все записи,вы можете вернуть только интересующие вас поля:

ksql> SELECT payload FROM TEST 
        WHERE PAYLOAD->source_order_id='3411976933216';
{MODIFIEDDATE=1552334325612, CREATEDDATE=1552334325612, CREATEDBY=C, SUCCESSFUL=true, SOURCE_ORDER_ID=3411976933216}

С KSQL вы можете записать результаты любого оператора SELECT в новую тему, которая заполняет его всеми существующими сообщениями вместе с каждым новым сообщением вИсходная тема фильтруется и обрабатывается в соответствии с объявленным оператором SELECT:

ksql> CREATE STREAM TEST_CREATED_BY_A AS
        SELECT * FROM TEST WHERE PAYLOAD->CREATEDBY='A';

Message
----------------------------
Stream created and running
----------------------------

Список тем в кластере Kafka:

ksql> SHOW TOPICS;

Kafka Topic            | Registered | Partitions | Partition Replicas | Consumers | ConsumerGroups
----------------------------------------------------------------------------------------------------
orders                 | true       | 1          | 1                  | 1         | 1
pageviews              | false      | 1          | 1                  | 0         | 0
products               | true       | 1          | 1                  | 1         | 1
TEST                   | true       | 1          | 1                  | 1         | 1
TEST_CREATED_BY_A      | true       | 4          | 1                  | 0         | 0

Распечатайте содержимое новой темы:

ksql> PRINT 'TEST_CREATED_BY_A' FROM BEGINNING;
Format:JSON
{"ROWTIME":1552475910106,"ROWKEY":"null","HEADER":{"SCHEMAVERSIONNO":"1"},"PAYLOAD":{"MODIFIEDDATE":1552334325212,"CREATEDDATE":1552334325212,"CREATEDBY":"A","SUCCESSFUL":true,"SOURCE_ORDER_ID":"3411976933214"}}
...