标签:Flink 1.10 SQL flink kafka connector sql org
一、简介
我这里使用的版本的是1.10.(之前试用过1.13的,结果可能是因为后面安装的Kafka插件版本和flink版本不同,导致查询出现问题)。安装这里就不说了。
1.启动flink
bin/start-cluster.sh。启动后可以看到如下进程。
2.启动sql-client
3.连接kafka需要增加扩展包(默认的只支持csv,file等文件系统)
flink-json-1.10.2.jar (https://repo1.maven.org/maven2/org/apache/flink/flink-json/1.10.2/flink-json-1.10.2.jar)
flink-sql-connector-kafka_2.11-1.10.2.jar (https://repo1.maven.org/maven2/org/apache/flink/flink-sql-connector-kafka_2.11/1.10.2/flink-sql-connector-kafka_2.11-1.10.2.jar)
将jar包flink 文件夹下的lib文件夹下
4.启动sql-client
bin/sql-client.sh embedded -1 lib/
5.创建Kafka Table
CREATE TABLE test(a int,
b int,
c int)with('connector.type' = 'kafka',
'connector.version' = 'universal',
'connector.properties.group.id' = 'g2.group1',
'connector.properties.bootstrap.servers' = '192.168.1.201:9092,192.168.1.202:9092',
'connector.properties.zookeeper.connect' = '192.168.1.201:2181',
'connector.topic' = 'customer_test',
'connector.startup-mode' = 'earliest-offset',
'format.type' = 'json');
6.插入数据
insert into test(a,b,c)
values(1,1,2),(2,10,2),(3,1,20);
在这里我一开始flink的版本是1.13和新增的kafka插件版本不对应。(我觉得是版本问题。哈哈。希望有大佬可以解决)
java.lang.ClassCastException: org.apache.calcite.sql.SqlBasicCall cannot be cast to org.apache.calcite.sql.SqlIdentifier
7.执行查询
select * from CustomerStatusChangedEvent ;
标签:Flink,1.10,SQL,flink,kafka,connector,sql,org 来源: https://blog.csdn.net/weixin_46014712/article/details/117409539
本站声明: 1. iCode9 技术分享网(下文简称本站)提供的所有内容,仅供技术学习、探讨和分享; 2. 关于本站的所有留言、评论、转载及引用,纯属内容发起人的个人观点,与本站观点和立场无关; 3. 关于本站的所有言论和文字,纯属内容发起人的个人观点,与本站观点和立场无关; 4. 本站文章均是网友提供,不完全保证技术分享内容的完整性、准确性、时效性、风险性和版权归属;如您发现该文章侵犯了您的权益,可联系我们第一时间进行删除; 5. 本站为非盈利性的个人网站,所有内容不会用来进行牟利,也不会利用任何形式的广告来间接获益,纯粹是为了广大技术爱好者提供技术内容和技术思想的分享性交流网站。