要将Flink中的数据写入到MySQL数据库中,可以使用Flink的JDBC Sink来实现。

首先,需要在Flink的依赖中添加相应的MySQL驱动依赖。例如,如果使用的是MySQL 8.x版本,可以添加以下依赖:

<dependency>
    <groupId>mysql</groupId>
    <artifactId>mysql-connector-java</artifactId>
    <version>8.0.26</version>
</dependency>

然后,在Flink的代码中,可以使用JDBC Sink将数据写入到MySQL中。以下是一个示例代码:

import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.connectors.jdbc.JdbcSink;
import org.apache.flink.streaming.connectors.jdbc.JdbcSinkBuilder;

import java.sql.PreparedStatement;

public class FlinkToMySQLExample {
    public static void main(String[] args) throws Exception {
        // 创建Flink的执行环境
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

        // 创建一个DataStream,例如从Kafka中读取数据
        DataStream<Tuple2<String, Integer>> stream = env.addSource(...);

        // 将DataStream中的数据转换为JDBC的写入格式
        DataStream<Tuple2<String, Integer>> jdbcStream = stream.map(new MapFunction<Tuple2<String, Integer>, Tuple2<String, Integer>>() {
            @Override
            public Tuple2<String, Integer> map(Tuple2<String, Integer> value) throws Exception {
                return value;
            }
        });

        // 创建JDBC Sink
        JdbcSink.sink(
                "INSERT INTO table_name (column1, column2) VALUES (?, ?)",
                (PreparedStatement ps, Tuple2<String, Integer> value) -> {
                    ps.setString(1, value.f0);
                    ps.setInt(2, value.f1);
                },
                JdbcSinkBuilder.JdbcConnectionOptions.builder()
                        .withUrl("jdbc:mysql://localhost:3306/database_name")
                        .withDriverName("com.mysql.cj.jdbc.Driver")
                        .withUsername("username")
                        .withPassword("password")
                        .build()
        );

        // 将数据写入到MySQL中
        jdbcStream.addSink(JdbcSink.sink(
                "INSERT INTO table_name (column1, column2) VALUES (?, ?)",
                (PreparedStatement ps, Tuple2<String, Integer> value) -> {
                    ps.setString(1, value.f0);
                    ps.setInt(2, value.f1);
                },
                JdbcSinkBuilder.JdbcConnectionOptions.builder()
                        .withUrl("jdbc:mysql://localhost:3306/database_name")
                        .withDriverName("com.mysql.cj.jdbc.Driver")
                        .withUsername("username")
                        .withPassword("password")
                        .build()
        ));

        // 执行Flink任务
        env.execute("Flink to MySQL Example");
    }
}

在上述代码中,需要根据实际情况替换以下部分:

  • table_name:要写入的MySQL表名
  • column1column2:要写入的MySQL表的列名
  • jdbc:mysql://localhost:3306/database_name:MySQL数据库的连接URL,其中localhost:3306是MySQL服务器的地址和端口,database_name是要连接的数据库名称
  • usernamepassword:MySQL数据库的用户名和密码

执行上述代码后,Flink会将DataStream中的数据写入到MySQL数据库中

flink更新写入mysql

原文地址: http://www.cveoy.top/t/topic/iVmV 著作权归作者所有。请勿转载和采集!

免费AI点我,无需注册和登录