【问题标题】:Read pipe separated values in ksql在 ksql 中读取管道分隔值
【发布时间】:2020-06-23 03:32:51
【问题描述】:

我正在研究 POC,我必须读取管道分隔值文件并将这些记录插入 ms sql server。 我正在使用 confluent 5.4.1 来使用 value_delimiter 创建流属性。但它的例外:Delimeter only supported with DELIMITED format

1.启动 Confluent(版本:5.4.1)::

[Dev root @ myip ~]
# confluent local start
    The local commands are intended for a single-node development environment
    only, NOT for production usage. https://docs.confluent.io/current/cli/index.html

Using CONFLUENT_CURRENT: /tmp/confluent.vHhSRAnj
Starting zookeeper
zookeeper is [UP]
Starting kafka
kafka is [UP]
Starting schema-registry
schema-registry is [UP]
Starting kafka-rest
kafka-rest is [UP]
Starting connect
connect is [UP]
Starting ksql-server
ksql-server is [UP]
Starting control-center
control-center is [UP]
[Dev root @ myip ~]
# jps
49923 KafkaRestMain
50099 ConnectDistributed
49301 QuorumPeerMain
50805 KsqlServerMain
49414 SupportedKafka
52103 Jps
51020 ControlCenter
1741
49646 SchemaRegistryMain
[Dev root @ myip ~]
#

2。创建主题:

[Dev root @ myip ~]
# kafka-topics --create --zookeeper localhost:2181 --replication-factor 1 --partitions 1 --topic SampleData
Created topic SampleData.

3.向 SampeData 主题提供管道分隔数据

[Dev root @ myip ~]
# kafka-console-producer --broker-list localhost:9092 --topic SampleData <<EOF
> this is col1|and now col2|and col 3 :)
> EOF
>>[Dev root @ myip ~]
#

4.启动 KSQL::

[Dev root @ myip ~]
# ksql

                  ===========================================
                  =        _  __ _____  ____  _             =
                  =       | |/ // ____|/ __ \| |            =
                  =       | ' /| (___ | |  | | |            =
                  =       |  <  \___ \| |  | | |            =
                  =       | . \ ____) | |__| | |____        =
                  =       |_|\_\_____/ \___\_\______|       =
                  =                                         =
                  =  Streaming SQL Engine for Apache Kafka® =
                  ===========================================

Copyright 2017-2019 Confluent Inc.

CLI v5.4.1, Server v5.4.1 located at http://localhost:8088

Having trouble? Type 'help' (case-insensitive) for a rundown of how things work!

5.为现有主题声明架构:SampleData

ksql> CREATE STREAM sample_delimited (
>       column1 varchar(1000),
>       column2 varchar(1000),
>       column3 varchar(1000))
>       WITH (KAFKA_TOPIC='SampleData', VALUE_FORMAT='DELIMITED', VALUE_DELIMITER='|');

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

6.验证数据到 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 sample_delimited emit changes limit 1;
+---------------------------+---------------------------+---------------------------+---------------------------+---------------------------+
|ROWTIME                    |ROWKEY                     |COLUMN1                    |COLUMN2                    |COLUMN3                    |
+---------------------------+---------------------------+---------------------------+---------------------------+---------------------------+
|1584339233947              |null                       |this is col1               |and now col2               |and col 3 :)               |
Limit Reached
Query terminated

7.编写一个新的 Kafka 主题:SampleDataAvro,将 sample_delimited 流中的所有数据序列化为 Avro 格式流

ksql> CREATE STREAM sample_avro WITH (KAFKA_TOPIC='SampleDataAvro', VALUE_FORMAT='AVRO') AS SELECT * FROM sample_delimited;
Delimeter only supported with DELIMITED format
ksql>

8.上面的行给出了异常:

Delimeter only supported with DELIMITED format

9.加载 ms sql kafka 连接配置

confluent local load test-sink -- -d ./etc/kafka-connect-jdbc/sink-quickstart-mssql.properties

【问题讨论】:

    标签: apache-kafka ksqldb


    【解决方案1】:

    您需要指定分隔符的唯一时间是当您定义从 source 主题读取的流时。

    这是我的工作示例:

    1. 用竖线分隔的数据填充主题:

      $ kafkacat -b localhost:9092 -t SampleData -P<<EOF
      this is col1|and now col2|and col 3 :)
      EOF
      
    2. 在其上声明一个流

      CREATE STREAM sample_delimited ( 
              column1 varchar(1000),
              column2 varchar(1000),
              column3 varchar(1000)) 
              WITH (KAFKA_TOPIC='SampleData', VALUE_FORMAT='DELIMITED', VALUE_DELIMITER='|');
      
    3. 查询流以确保其正常工作

      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 sample_delimited emit changes limit 1;
      +----------------+--------+---------------+--------------+--------------+
      |ROWTIME         |ROWKEY  |COLUMN1        |COLUMN2       |COLUMN3       |
      +----------------+--------+---------------+--------------+--------------+
      |1583933584658   |null    |this is col1   |and now col2  |and col 3 :)  |
      Limit Reached
      Query terminated
      
    4. 将数据重新序列化到 Avro:

      CREATE STREAM sample_avro WITH (KAFKA_TOPIC='SampleDataAvro', VALUE_FORMAT='AVRO') AS SELECT * FROM sample_delimited;
      
    5. 转储主题的内容 - 注意现在是 Avro:

      ksql> print SampleDataAvro;
      Key format: UNDEFINED
      Value format: AVRO
      rowtime: 3/11/20 1:33:04 PM UTC, key: <null>, value: {"COLUMN1": "this is col1", "COLUMN2": "and now col2", "COLUMN3": "and col 3 :)"}
      

    您遇到的错误是由错误#4200 引起的。您可以等待 Confluent Platform 的下一个版本,或使用已修复问题的standalone ksqlDB

    这里使用 ksqlDB 0.7.1 将数据流式传输到 MS SQL:

    CREATE SINK CONNECTOR SINK_MSSQL WITH (
        'connector.class'     = 'io.confluent.connect.jdbc.JdbcSinkConnector',
        'connection.url'      = 'jdbc:sqlserver://mssql:1433',
        'connection.user'     = 'sa',
        'connection.password' = 'Admin123',
        'topics'              = 'SampleDataAvro',
        'key.converter'       = 'org.apache.kafka.connect.storage.StringConverter',
        'auto.create'         = 'true',
        'insert.mode'         = 'insert'
      );
    

    现在在 MS SQL 中查询数据

    1> Select @@version
    2> go
    
    ---------------------------------------------------------------------
    Microsoft SQL Server 2017 (RTM-CU17) (KB4515579) - 14.0.3238.1 (X64)
            Sep 13 2019 15:49:57
            Copyright (C) 2017 Microsoft Corporation
            Developer Edition (64-bit) on Linux (Ubuntu 16.04.6 LTS)
    
    (1 rows affected)
    
    1> SELECT * FROM SampleDataAvro;
    2> GO
    COLUMN3        COLUMN2         COLUMN1     
    -------------- --------------- ------------------
    and col 3 :)   and now col2    this is col1
    
    (1 rows affected)
    

    【讨论】:

    • 我仍然收到错误ksql&gt; CREATE STREAM sample_avro WITH (KAFKA_TOPIC='SampleDataAvro', VALUE_FORMAT='AVRO') AS SELECT * FROM sample_delimited; Delimeter only supported with DELIMITED format
    • 我正在使用,Copyright 2017-2019 Confluent Inc. CLI v5.4.1, Server v5.4.1 located at localhost:8088
    • 这与您的问题不同。您可以编辑您的问题并将您的会话的完整输出附加到它,显示create stream 和随后的create stream as select 有错误?
    • 我已经更新了我的问题,在控制台上提供了完整的输出。基本上我 试图读取管道分隔值文件并沉入 ms sql 数据库。
    • 我已经修改了答案
    猜你喜欢
    • 1970-01-01
    • 2020-08-14
    • 2020-12-15
    • 1970-01-01
    • 2011-10-04
    • 1970-01-01
    • 2018-05-16
    • 1970-01-01
    • 2021-05-24
    相关资源
    最近更新 更多