flink消费kafka数据的版本问题,可以去https://mvnrepository.com/,查看对应版本。
如果在开发过程中,出现版本不对应,那么kafka的topic一定要重新创建一个,以防各种错误。
环境:
mysql
zookeeper:3.4.13
kafka:0.8_2.11
flink:1.7.2(pom.xml中)
启动zookeeper
bin/zkServer.sh start
启动kafka:(此时我的topic是lhqtest)
完整代码:
pom.xml:
代码FlinkKafkajson:
import org.apache.flink.api.common.typeinfo.Types;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.api.TableEnvironment;
import org.apache.flink.table.api.java.StreamTableEnvironment;
import org.apache.flink.table.descriptors.Json;
import org.apache.flink.table.descriptors.*;
import org.apache.flink.table.descriptors.Schema;
import org.apache.flink.types.Row;
/**
* Created by luhaiqing on 2019/6/11.
*/
public class FlinkKafkajson {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnvironment = TableEnvironment.getTableEnvironment(env);
tableEnvironment.connect(new Kafka().version("0.8").topic("lhqtest").startFromLatest()
.property("bootstrap.servers","192.168.190.133:9092")
.property("zookeeper.connect","192.168.190.133:2181")
.property("group.id", "lhqtest"))
.withFormat(new Json().failOnMissingField(true).deriveSchema())
.withSchema(new Schema()
.field("id", Types.INT)
.field("name", Types.STRING)
.field("sex", Types.STRING)
)
.inAppendMode()
.registerTableSource("lhq_user");
Table table = tableEnvironment.scan("lhq_user").select("id,name,sex");
DataStream<Row> personDataStream = tableEnvironment.toAppendStream(table,Row.class);
personDataStream.addSink(new MysqlSink());
env.execute("userPv from Kafka");
}
}
写入mysql代码:
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;
import org.apache.flink.types.Row;
import java.sql.Connection;
import java.sql.DriverManager;
import java.sql.PreparedStatement;
/**
* Created by luhaiqing on 2019/6/5.
*/
public class MysqlSink extends RichSinkFunction<Row>
{
private Connection connection;
private PreparedStatement preparedStatement;
@Override
public void open(Configuration parameters) throws Exception {
super.open(parameters);
String className = "com.mysql.jdbc.Driver";
Class.forName(className);
String url = "jdbc:mysql://localhost:3306/test";
String user = "root";
String password = "123456";
connection = DriverManager.getConnection(url, user, password);
String sql = "replace into flinkjsontest(id,name,sex) values(?,?,?)";
preparedStatement = connection.prepareStatement(sql);
super.open(parameters);
}
@Override
public void close() throws Exception {
super.close();
if (preparedStatement != null) {
preparedStatement.close();
}
if (connection != null) {
connection.close();
}
super.close();
}
public void invoke(Row value, Context context) throws Exception {
int id = (int)value.getField(0);
String name = (String)value.getField(1);
String sex = (String)value.getField(2);
System.out.print(id+":"+name+":"+sex);
preparedStatement.setInt(1, id);
preparedStatement.setString(2, name);
preparedStatement.setString(3,sex);
int i = preparedStatement.executeUpdate();
if (i > 0) {
System.out.println("value=" + value);
}
}
}
看到这里,给你赞吧!!!
原网址: 访问
创建于: 2020-05-13 05:00:40
目录: default
标签: 无
未标明原创文章均为采集,版权归作者所有,转载无需和我联系,请注明原出处,南摩阿彌陀佛,知识,不只知道,要得到
java windows火焰图_mob64ca12ec8020的技术博客_51CTO博客 - 在windows下不可行,不知道作者是怎样搞的 监听SpringBoot 服务启动成功事件并打印信息_监听springboot启动完毕-CSDN博客 SpringBoot中就绪探针和存活探针_management.endpoint.health.probes.enabled-CSDN博客 u2u转换板 - 嘉立创EDA开源硬件平台 Spring Boot 项目的轻量级 HTTP 客户端 retrofit 框架,快来试试它!_Java精选-CSDN博客 手把手教你打造一套最牛的知识笔记管理系统! - 知乎 - 想法有重合-理论可参考 安宇雨 闲鱼 机械键盘 客制化 开贴记录 文本 linux 使用find命令查找包含某字符串的文件_beijihukk的博客-CSDN博客_find 查找字符串 ---- mac 也适用 安宇雨 打字音 记录集合 B站 bilibili 自行搭建 开坑 真正的客制化 安宇雨 黑苹果开坑 查找工具包maven pom 引用地 工具网站 Dantelis 介绍的玩轴入坑攻略 --- 关于轴的一些说法 --- 非官方 ---- 心得而已 --- 长期开坑更新 [本人问题][新开坑位]关于自动化测试的工具与平台应用 机械键盘 开团 网站记录 -- 能做一个收集的程序就好了 不过现在没时间 -- 信息大多是在群里发的 - 你要让垃圾佬 都去一个地方看难度也是很大的 精神支柱 [超级前台]sprinbboot maven superdesk-app 记录 [信息有用] [环境准备] [基本完成] [sebp/elk] 给已创建的Docker容器增加新的端口映射 - qq_30599553的博客 - CSDN博客 [正在研究] Elasticsearch, Logstash, Kibana (ELK) Docker image documentation elasticsearch centos 安装记录 及 启动手记 正式服务器 39 elasticsearch 问题合集 不断更新 6.1.1 | 6.5.1 两个版本 博客程序 - 测试 - bug记录 等等问题 laravel的启动过程解析 - lpfuture - 博客园 OAuth2 Server PHP 用 Laravel 搭建带 OAuth2 验证的 RESTful 服务 | Laravel China 社区 - 高品质的 Laravel 和 PHP 开发者社区 利用Laravel 搭建oauth2 API接口 附 Unauthenticated 解决办法 - 煮茶的博客 - SegmentFault 思否 使用 OAuth2-Server-php 搭建 OAuth2 Server - 午时的海 - 博客园 基于PHP构建OAuth 2.0 服务端 认证平台 - Endv - 博客园 Laravel 的 Artisan 命令行工具 Laravel 的文件系统和云存储功能集成 浅谈Chromium中的设计模式--终--Observer模式 浅谈Chromium中的设计模式--二--pre/post和Delegate模式 浅谈Chromium中的设计模式--一--Chromium中模块分层和进程模型 DeepMind 4 Hacking Yourself README.md update 20211011
Laravel China 简书 知乎 博客园 CSDN博客 开源中国 Go Further Ryan是菜鸟 | LNMP技术栈笔记 云栖社区-阿里云 Netflix技术博客 Techie Delight Linkedin技术博客 Dropbox技术博客 Facebook技术博客 淘宝中间件团队 美团技术博客 360技术博客 古巷博客 - 一个专注于分享的不正常博客 软件测试知识传播 - 测试窝 有赞技术团队 阮一峰 语雀 静觅丨崔庆才的个人博客 软件测试从业者综合能力提升 - isTester IBM Java 开发 使用开放 Java 生态系统开发现代应用程序 pengdai 一个强大的博主 HTML5资源教程 | 分享HTML5开发资源和开发教程 蘑菇博客 - 专注于技术分享的博客平台 个人博客-leapMie 流星007 CSDN博客 - 舍其小伙伴 稀土掘金 Go 技术论坛 | Golang / Go 语言中国知识社区
最新评论