一、jdbc连接
1、引入依赖
<!-- https:
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-connector-jdbc_2.12</artifactId>
<version>1.14.4</version>
<scope>provided</scope>
</dependency>
2、实体类
public class UserInfo {
private Integer userId;
private String userName;
private String userRealName;
private String userPwd;
private String userTel;
private String userEmail;
private Integer userStatus;
private Date userCreateTime;
private Date userUpdateTime;
}
二、自定义 SourceFunction
public class MysqlRichParallelSource extends RichParallelSourceFunction<User> {
private boolean close = false;
@Override
public void run(SourceContext<User> out) throws Exception {
String url = "jdbc:mysql://192.168.100.88:3306/newframe?useUnicode=true&characterEncoding=UTF-8&serverTimezone=Asia/Shanghai&useSSL=false&allowMultiQueries=true&rewriteBatchedStatements=true";
String sql = "select * from user_info";
Connection conn = null;
PreparedStatement ps = null;
ResultSet rs = null;
try {
conn = DriverManager.getConnection(url, "root", "root");
ps = conn.prepareStatement(sql);
} catch (SQLException e) {
e.printStackTrace();
}
while (!close) {
rs = ps.executeQuery();
while (rs.next()) {
Integer userId = rs.getInt("user_id");
String userName = rs.getString("user_name");
String userRealName = rs.getString("user_real_name");
User user = new User()
.setUserId(userId)
.setUserName(userName)
.setUserRealName(userRealName);
out.collect(user);
}
Thread.sleep(5000);
cancel();
}
close(conn, ps, rs);
}
@Override
public void cancel() {
close = true;
}
public void close(Connection conn, PreparedStatement ps, ResultSet rs) throws Exception {
if (conn != null) {
conn.close();
}
if (ps != null) {
ps.close();
}
if (rs != null) {
rs.close();
}
}
}
三、测试
@Test
public void fromMysqlTest() throws Exception {
final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setRuntimeMode(RuntimeExecutionMode.STREAMING);
DataStreamSource<User> source = env.addSource(new MysqlRichParallelSource());
source.print();
env.execute("flink streaming from mysql");
}
结果:多线程同时从mysql获取数据,启动了8个线程,每个数据获取了8次
2> User(userId=20, userName=wangwu2, userRealName=王五2, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
8> User(userId=20, userName=wangwu2, userRealName=王五2, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
6> User(userId=20, userName=wangwu2, userRealName=王五2, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
3> User(userId=20, userName=wangwu2, userRealName=王五2, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
5> User(userId=20, userName=wangwu2, userRealName=王五2, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
4> User(userId=20, userName=wangwu2, userRealName=王五2, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
1> User(userId=20, userName=wangwu2, userRealName=王五2, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
7> User(userId=20, userName=wangwu2, userRealName=王五2, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
4> User(userId=22, userName=user1, userRealName=用户1, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
5> User(userId=22, userName=user1, userRealName=用户1, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
3> User(userId=22, userName=user1, userRealName=用户1, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
6> User(userId=22, userName=user1, userRealName=用户1, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
6> User(userId=23, userName=user2, userRealName=用户2, userPwd=null, userTel=null, userEmail=null, userStatus=null, userCreateTime=null, userUpdateTime=null)
|