Flink教程(5)-Flink常用API

发布时间:2026/7/21 12:53:25
Flink教程(5)-Flink常用API
Flink教程(5)-Flink常用API【订阅专栏合集,作者所有付费文章都能看】文章目录FlinkAPIEnvironmentSourceTransformationFlink数据类型SinkFlinkAPIEnvironment执行Flink程序首先要判断flink环境。Flink中有3种获取执行环境的方式。1getExecutionEnvironment获取当前执行程序的上下文。如果是直接在IDEA中运行的JAVA代码则此方法返回本地执行环境。如果是从命令行或web页面提交flink任务到集群中则此方法返回的是集群执行环境。这种方式是最常用的Flink底层帮我们判断具体调用本地还是远程环境。ExecutionEnvironment.getExecutionEnvironment();//获取批处理执行环境 StreamExecutionEnvironment.getExecutionEnvironment(); //获取流处理执行环境复制2createLocalEnvironment直接返回本地执行环境这种方式可以指定并行度。如不指定则使用当前机器可用cpu核数作为并行度。其实第1种方式判断当前环境是本地环境的话底层也会调此方法。ExecutionEnvironment.createLocalEnvironment();复制3createRemoteEnvironment获取远程集群执行环境。如果将Jar包提交到远程Flink集群执行则需指定JobManager的IP和port并指定jar包路径ExecutionEnvironment.createRemoteEnvironment(hostname,port,hdfs://wordCount.jar);复制Sourcesource是Flink应用程序的数据来源。作为一款通用的数据处理框架flink既可以处理静态的历史数据集也可以处理实时的流式数据。流式计算场景下只要数据源源不断传入flink就能一直处理。下面讲解Flink中的几种数据输入方式。1从本地集合中读取executionEnvironment.fromCollection(Arrays.asList(a, b, c,d));//从JAVA Collection中读取数据 executionEnvironment.fromElements(1, 2, 3, 4);//从给定的对象序列中读取数据复制2从文件中读取String inputPath F:\\data\\file; executionEnvironment.readTextFile(inputPath);//使用默认的文件格式复制3从socket中读取env.socketTextStream(localhost, 9999);//从指定的IP地址和端口处读取数据使用默认行分隔符 env.socketTextStream(hostname, port, delimiter);//指定行分隔符复制4从Kafka中读取实际开发中Kafka作为Flink数据源非常常见可以说Kafka和Flink在流式数据处理领域是天生的一对。引入Kafka连接器pom依赖连接器的版本和Flink版本保持一致dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-kafka-0.11_2.11/artifactId version1.9.2/version /dependency复制Flink中添加kafka数据源public static void main(String[] args) throws Exception { final StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); //kafka配置参数 Properties props new Properties(); props.put(bootstrap.servers, 192.168.174.129:9092); props.put(zookeeper.connect, 192.168.174.129:2181); //props.put(group.id, metric-group); props.put(key.deserializer, org.apache.kafka.common.serialization.StringDeserializer); props.put(value.deserializer, org.apache.kafka.common.serialization.StringDeserializer); props.put(auto.offset.reset, latest); //FlinkKafkaConsumer011 表示对应的kafka版本是0.11.x DataStreamSourceString dataStreamSource env.addSource(new FlinkKafkaConsumer011( test01, //kafka topic new SimpleStringSchema(), // String 序列化 props)).setParallelism(1); dataStreamSource.print(); //把从 kafka 读取到的数据打印在控制台 env.execute(Flink add kafka data source); }复制addSource是一般化的添加数据源的算子前面几种source都是Flink根据特定应用场景封装好的算子底层还是调用了addSource。可以测试一下上述程序。在linux服务器上启动kafka集群并通过命令行运行一个Producer发送消息。查看Flink是否消费到数据。Kafka相关教程可以参考这篇文章《Kafka 实战教程》5自定义source有时为了方便测试Flink应用程序我们需要手动造数据这就要用到自定义数据源。自定义的DataSource只要实现org.apache.flink.streaming.api.functions.source.SourceFunction接口即可被作为数据源添加。下面的例子展示了如何自定义数据源。需求是实现一个实时数字生成器1秒钟产生1个自增数字发送到Flink。Flink收到数据后放大两倍输出。/** * 自定义Flink数据源,重写SourceFunction的run和cancel方法 */ import org.apache.flink.streaming.api.functions.source.SourceFunction; public class MyDataSource implements SourceFunctionInteger { private boolean isRunning true; /** * run方法里编写数据产生逻辑 * param ctx * throws Exception */ Override public void run(SourceContextInteger ctx) throws Exception { int i 1; while (isRunning) { ctx.collect(i); i; Thread.sleep(1000); } } Override public void cancel() { isRunning false; } }复制public class MyDataSourceTest { public static void main(String[] args) throws Exception { final StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); DataStreamSourceInteger mySource env.addSource(new MyDataSource()); SingleOutputStreamOperatorInteger res mySource.map(e - 2 * e); res.print(); env.execute(Flink add dataSource); } }复制掌握了自定义数据源的使用有助于实际开发比如可以写个实时读取MySql数据源的工具TransformationTransform也可称为Operator翻译为中文为“算子”其实就是数据转换操作。下面讲解Flink中的几种数据转换操作先从流式处理即DataStream 操作讲起。批处理与之类似。1MapMap就是映射顾名思义就是将输入数据进行转换操作。map算子的输入参数是一个MapFunction我们只要实现它重写其中的map函数即可public interface MapFunctionT, O extends Function, Serializable { O map(T value) throws Exception; }复制比如将商品数据流中的每个商品价格翻倍SingleOutputStreamOperatorProduct map dataStreamSource.map(new MapFunctionProduct, Product() { Override public Product map(Product product) throws Exception { product.price product.price * 2; return product; } }); map.print();复制对于简单的转换操作我们也可以直接使用lambda 表达式比如.map(e - 2 * e);StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); DataStreamSourceInteger mySource env.addSource(new MyDataSource()); SingleOutputStreamOperatorInteger res mySource.map(e - 2 * e);//把数据源中的每个元素放大2倍 res.print();复制2FlatMapFlatMap意指扁平化的map即将每个元素map后的数据打散重新组成一个“宽”的集合。和JDK8中的flatMap本质一样。FlatMap的输入参数是一个FlatMapFunction只要重写其flatMap方法即可value是输入数据out是输出数据收集器public interface FlatMapFunctionT, O extends Function, Serializable { void flatMap(T value, CollectorO out) throws Exception; }复制dataStreamSource.flatMap(new FlatMapFunctionString, Tuple2String, Integer() { //接收一个字符串Wordcount中表示一行数据输出一个2元组 Override public void flatMap(String value, CollectorTuple2String, Integer out) throws Exception { String[] splits value.split(\\s); for (String word : splits) { out.collect(new Tuple2String, Integer(word, 1)); } } })复制FlatMap和Map的区别在第一篇快速入门案例中已讲解此处不再赘述。3Filter对元素进行过滤重写FilterFunction的filter实现过滤逻辑简单的过滤逻辑可以直接使用lambda表达式public interface FilterFunctionT extends Function, Serializable { boolean filter(T value) throws Exception; }复制例如过滤价格超过100的商品SingleOutputStreamOperatorProduct res mySource.filter(new FilterFunctionProduct() { Override public boolean filter(Product product) throws Exception { if (product.price 100) { return true; } return false; } }) res.print();复制4KeyBy根据指定的key对流数据元素进行分区底层基于hash算法hashCode相同的key被分到同一个分区即分到下游算子并行节点中的一个。比如快速入门案例中flatMap之后的数据按照单词分组即按照二元组数据的第一个字段Tuple2.f0 分组DataStreamTuple2String, Integer dataStream dataStreamSource .flatMap(new Splitter()) .keyBy(value - value.f0)复制keyBy的参数是KeySelectorIN, KEY前一个泛型表示来源数据类型后一个泛型表示从原数据中提取出来的key的类型再比如根据商品的品牌来分组KeyedStreamProduct, String keyByedProd productStream.keyBy(new KeySelectorProduct, String() { Override public String getKey(Product product) throws Exception { return product.brand; } }); keyByedProd.print();复制简写.keyBy(product- product.brand)5Reducereduce俗称“约减”就是将元素进行聚合处理。常见的sum、min、max、count、average等聚合操作都可以使用原生的reduce实现。reduce算子的入参是ReduceFunctionvalue1表示前一个元素value2表示后一个元素reduce方法是具体的数据处理逻辑。reduce操作实质上就是不断地将数据源中两个值合并为同一类型的一个值reduce函数连续应用于输入数据流中的所有值直到只剩下一个值聚合之后的结果。public interface ReduceFunctionT extends Function, Serializable { T reduce(T value1, T value2) throws Exception; }复制比如统计各个品牌商品的总价SingleOutputStreamOperatorProduct reduceRes productStream.keyBy(new KeySelectorProduct, String() { Override public String getKey(Product product) throws Exception { return product.brand; } }).reduce(new ReduceFunctionProduct() { Override public Product reduce(Product product1, Product product2) throws Exception { product2.price (product1.price product2.price); return product2; } }); reduceRes.print();复制6AggregationFlink中支持对数据流的各种聚合操作并封装了很多聚集函数。像min、max、sum等聚集函数都可以应用于 KeyedStream获得聚合结果。聚合算子参数如果是int类型则表示聚合字段的下标(从0开始)。比如快速入门案例中对二元组数据求和.sum(1)表示求二元组中第二个字段单词计数的和。如果是string类型则表示聚合字段名通常是一个pojo对象的public属性。KeyedStream.sum(0) KeyedStream.sum(field0) KeyedStream.min(1) KeyedStream.min(field1) KeyedStream.max(2) KeyedStream.max(field2) KeyedStream.minBy(3) KeyedStream.minBy(field3) KeyedStream.maxBy(4) KeyedStream.maxBy(field4)复制7Split和SelectSplit是根据指定条件将数据流拆分为两个或多个流可以单独处理每个数据流。Select是从拆分的流中选择特定的流。select和split一般结合使用正如keyBy和聚集函数一起使用一样。实现这样的需求按照商品价格比如100元为界将商品分为优品(100)和良品。public static void main(String[] args) throws Exception { StreamExecutionEnvironment executionEnvironment StreamExecutionEnvironment.getExecutionEnvironment(); ListProduct products new ArrayList(); products.add(new Product(A, 阿迪, 990)); products.add(new Product(B, 安踏, 90)); products.add(new Product(C, 耐克, 880)); products.add(new Product(D, 特步, 80)); DataStreamSourceProduct streamSource executionEnvironment.fromCollection(products); SplitStreamProduct splitStream streamSource.split(new OutputSelectorProduct() { Override public IterableString select(Product product) { ListString list new ArrayList();//使用list作为临时数据结构存储标签 if (product.getPrice() 100) { list.add(优品); } else { list.add(良品); } return list; } }); DataStreamProduct superiorProducts splitStream.select(优品); DataStreamProduct acceptedProducts splitStream.select(良品); DataStreamProduct allProducts splitStream.select(良品,优品); superiorProducts.print(优品); //启动计算任务 executionEnvironment.execute(Stream operator); }复制控制台输出“优品”的数据优品:1 Product{nameC, brand耐克, price880.0} 优品:4 Product{nameA, brand阿迪, price990.0}复制有分流操作那么与之对应的必然有合流操作。Flink中合流操作有2种Union和Connect。8UnionUnion函数表示将两个或多个数据类型相同的流组合在一起即求并集。StreamExecutionEnvironment executionEnvironment StreamExecutionEnvironment.getExecutionEnvironment(); DataStreamSourceString dataStream1 executionEnvironment.fromCollection(Arrays.asList(a, b, c, d)); DataStreamSourceString dataStream2 executionEnvironment.fromCollection(Arrays.asList(1, 2, 3, 4)); DataStreamString union dataStream2.union(dataStream1); union.print();复制控制台输出合并之后的数据3 b 4 c 1 d 1 3 4 2 2 a 3 1 2 4复制9connect和CoMap两个datastream连接后转变成connectedstreams即datastream,datastream-connectedstreams。与union不同的是connect不要求被连接的两个流数据类型相同。两个流虽然被connect到了同一个流中但是合并之后的流内部依然保持各自的数据格式不变相互独立。connect通常和coMap一起使用coMap对connect之后的流做数据处理。实际应用中一个数据流过来可能先根据元素的某种特征分开处理到了一定阶段又需要合并处理此时就需要用到分流和合流操作。实现这样的需求连接一个二元组类型的数据流和一个Product的数据流。并分别对连接后的数据流做map操作。还是沿用之前的例子:public static void main(String[] args) throws Exception { StreamExecutionEnvironment executionEnvironment StreamExecutionEnvironment.getExecutionEnvironment(); ListProduct products new ArrayList(); products.add(new Product(A, 阿迪, 999)); products.add(new Product(B, 安踏, 99)); products.add(new Product(C, 耐克, 888)); products.add(new Product(D, 特步, 89)); DataStreamSourceProduct streamSource executionEnvironment.fromCollection(products); SplitStreamProduct splitStream streamSource.split(new OutputSelectorProduct() { Override public IterableString select(Product product) { ListString list new ArrayList(); if (product.getPrice() 100) { list.add(优品); } else { list.add(良品); } return list; } }); DataStreamProduct superiorProducts splitStream.select(优品); DataStreamProduct acceptedProducts splitStream.select(良品); //先将优品数据转换为2元组类型 DataStreamTuple2String, Double superiorProductsStream superiorProducts.map(new MapFunctionProduct, Tuple2String, Double() { Override public Tuple2String, Double map(Product product) throws Exception { return new Tuple2(product.getName(), product.getPrice()); } }); //连接二元组数据流和Product数据类型并做算子操作都转换为3元组 SingleOutputStreamOperatorObject operator superiorProductsStream.connect(acceptedProducts).map(new CoMapFunctionTuple2String, Double, Product, Object() { Override public Object map1(Tuple2 value) throws Exception { return new Tuple3(value.f0, value.f1, 优品); } Override public Object map2(Product value) throws Exception { return new Tuple3(value.getName(), value.getPrice(), 良品); } }); operator.print(coMap); //启动计算任务 executionEnvironment.execute(Stream operator); }复制控制台输出如下内容说明案例中不同数据类型的流连接(connect)、处理(coMap)成功。coMap:1 (A,999.0,优品) coMap:2 (C,888.0,优品) coMap:4 (D,89.0,良品) coMap:3 (B,99.0,良品)复制观察上面的介绍的几种操作可以总结一些规律。比如keyBy操作总是和聚集函数一起使用、split通常和select一起使用、connect和coMap一起使用。datastream split后得到splitstream再select之后又转换为datastream 同样的datastream connect之后得到connectedstreams再经coMap操作后又转换为datastream 。union合并的两个流数据类型必须相同合并过程不涉及流类型的转换。而connect不要求数据流的元素类型相同。union操作可以操作多个流connect操作只能操作两个流。上述介绍的 DataStream流处理 数据转换操作中有些也适合DataSet批处理。比如 Map、FlatMap、Reduce、Filter 等。当然DataSet也有一些特有算子。比如在DataStream中分区是 KeyBy而DataSet中是GroupBy。这在快速入门案例中已经演示不再赘述。DataSet有个first(n)方法可以返回DataSet中前 n个元素比如env.readTextFile(inputPath).first(2);返回数据集中前2个元素。Flink数据类型前文中介绍Flink中算子的使用时提到了数据类型下面简单介绍一下Flink中所支持的数据类型。Flink应用程序处理的是由数据对象组成的连续不断的数据流。这些数据对象需要被序列化和反序列化以便能够通过网络传输以及从检查点、保存点、状态后端存储读取。为了明确应用程序所处理的数据类型Flink底层提供了一套完备的数据类型信息并且为每一种类型提供了序列化器、反序列化器以及比较器。此外Flink还提供了类型提取系统自动分析函数的输入类型和输出类型以获得对应的序列化器和反序列化器。在使用lambda函数或者泛型类型时需显式指定类型信息。Flink DataStream里的元素类型支持JAVA和Scala中的所有基本类型像Int、Long、Double、String等。此外还支持Tuple元组类型、Java简单对象(pojo)、scala样例类以及一些集合类型比如Java的ArrayList、HashMap、Enum等。Flink的每个函数都提供了对应的Rich版本。富函数相比普通的函数可以获取flink运行时上下文、生命周期方法。生命周期方法中通常可以做一些初始化及收尾操作比如连接数据库、关闭数据库连接。Sinksink顾名思义下沉在Flink中意指数据输出、数据落地的意思。最简单的数据输出方式就是打印到控制台调用datastream的print()方法即可print就是一种sink操作。对于不同的sink方式Flink提供了各种内置的输出格式。除了基本的输入输出数据源外flink目前还支持下列第三方组件作为数据源。Apache Kafka(source/sink)Apache Cassandra(sink)Amazon Kinesis Streams(source/sink)Elasticsearch(sink)Hadoop FileSystem(sink)RabbitMQ(source/sink)Apache NiFi(source/sink)Twitter Streaming API(source)Google PubSub(source/sink)JDBC(sink)本节介绍几种常用的数据输出方式。1普通文件、socketwriteAsText()/TextOutputFormat将元素按行写入字符串。字符串通过调用每个元素的toString()方法获得。writeAsCsv(…)/CsvOutputFormat将数据以逗号分隔的形式写入文件。换行符和字段分隔符可配置。每个字段的值来自对象的toString()方法。writeUsingOutputFormat() / FileOutputFormat自定义文件输出格式支持自定义对象到字节的转换。writeToSocket根据指定格式(Serialization Schema)将元素写入网络套接字。举例如下public class SinkDemo { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); //设置并行度为1输出结果全部写出到一个文件否则开发环境会使用默认并行度分区分区数为当前机器逻辑cpu核数 env.setParallelism(1); //准备数据源 ListTuple2String, Integer list new ArrayList(); list.add(new Tuple2(A, 100)); list.add(new Tuple2(B, 200)); list.add(new Tuple2(C, 300)); list.add(new Tuple2(D, 400)); DataStreamSourceTuple2String, Integer dataStreamSource env.fromCollection(list); dataStreamSource.print(); //除了路径参数是必填外还可以通过指定第二个参数来定义输出模式 dataStreamSource.writeAsText(d://sink-text.txt, FileSystem.WriteMode.OVERWRITE); //如果想要将输出结果全部写出到一个文件可以单独设置算子的并行度为 1 dataStreamSource.writeAsCsv(d://sink-csv.txt, FileSystem.WriteMode.OVERWRITE, \n, ,).setParallelism(1); //自定义的输出格式writeAsText/writeAsCsv底层调用的都是该方法 dataStreamSource.writeUsingOutputFormat(new TextOutputFormat(new Path(d://sink-file.txt), UTF-8)); //以字符串的形式输出到socket服务器 dataStreamSource.map(t - t.f0 : t.f1 \r\n).writeToSocket(192.168.244.131, 9999, new SimpleStringSchema()); env.execute(sink demo); } }复制测试socket输出时先在linux服务器上使用nc -lk 9999 模拟socket服务器开启监听。2kafkakafka和flink天生对流式数据友好因此实际生产中经常搭配使用。比如flink从数据源接收到数据处理完成后再发送一个消息到kafka中任其消费。也有从kafka进、kafka出的使用场景即输入、输出源都是kafka。比如对原始输出数据进行分流处理并且处理完成后发送到不同的消费者topic中去。下面介绍如何在flink中集成kafka。引入依赖dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-kafka-0.11_2.11/artifactId version1.11.0/version /dependency复制需求实现Flink消费kafka消息队列的消息转换处理后再次输出到kafka中。实体类Productpublic class Product { private String name; private String brand; private double price; //省略get/set }复制通过linux命令行创建2个topic一个由flink消费另一个由flink写入。bin/kafka-topics.sh --create \ --bootstrap-server 192.168.244.131:9092 \ --replication-factor 1 \ --partitions 1 \ --topic flink-stream-in-topic复制bin/kafka-topics.sh --create \ --bootstrap-server 192.168.244.131:9092 \ --replication-factor 1 \ --partitions 1 \ --topic flink-stream-out-topic复制查看topicbin/kafka-topics.sh --list --bootstrap-server 192.168.244.131:9092启动一个消费者接收flink的输出bin/kafka-console-consumer.sh --bootstrap-server 192.168.244.131:9092 --topic flink-stream-out-topic启动一个生产者向flink应用程序监听的topic发送消息bin/kafka-console-producer.sh --topic flink-stream-in-topic --bootstrap-server 192.168.244.131:9092在生产者端输入json串:{name:跑鞋,brand:Nike,price:1000}Flink应用程序集成Kafkapublic class KafkaSink { public static void main(String[] args) throws Exception { final StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); //kafka配置 Properties props new Properties(); props.put(bootstrap.servers, 192.168.244.131:9092); props.put(zookeeper.connect, 192.168.244.131:2181); props.put(key.deserializer, org.apache.kafka.common.serialization.StringDeserializer); //key 反序列化 props.put(value.deserializer, org.apache.kafka.common.serialization.StringDeserializer);//value 反序列化 props.put(auto.offset.reset, latest); //从kafka中消费数据 DataStreamSourceString dataStreamSource env.addSource(new FlinkKafkaConsumer011( flink-stream-in-topic, //kafka topic new SimpleStringSchema(), // String序列化 props)).setParallelism(1);//并行度一般不超过kafka topic分区数 dataStreamSource.print(); //把从 kafka 读取到的数据打印在控制台 //对数据进行业务处理 SingleOutputStreamOperatorString streamOperated dataStreamSource.map(new MapFunctionString, String() { Override public String map(String value) throws Exception { Product product JSON.parseObject(value, Product.class); //将商品价格翻倍 product.setPrice(product.getPrice()*2); return JSON.toJSONString(product); } }); //将处理完的数据再次发送到kafka中 streamOperated.addSink(new FlinkKafkaProducer011( flink-stream-out-topic, new SimpleStringSchema(), props)).setParallelism(1); env.execute(kafka data source); } }复制运行flink应用后kafka消费者端将收到处理后的数据{brand:Nike,name:跑鞋,price:2000.0}3redisRedis Connector 用于向 Redis 发送数据。可以使用三种不同的方法与不同类型的 Redis 环境进行通信单 Redis 服务器Redis 集群Redis Sentinel(哨兵)不同模式主要是Config类的不同本例展示了单机模式下Flink写入redis引入redis连接器依赖dependency groupIdorg.apache.bahir/groupId artifactIdflink-connector-redis_2.11/artifactId version1.0/version /dependency复制编写flink应用代码public class RedisSinkDemo { public static void main(String[] args) throws Exception { final StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); //准备数据源,简单起见,这里使用本地集合数据 ListTuple2String, String list new ArrayList(); list.add(new Tuple2(A, apple)); list.add(new Tuple2(B, bird)); list.add(new Tuple2(C, cat)); list.add(new Tuple2(D, dog)); DataStreamSourceTuple2String, String dataStreamSource env.fromCollection(list); //单机Redis配置,这里只简单配置ip/端口,还支持其它配置比如maxTotal、maxIdle、timeout FlinkJedisPoolConfig redisConf new FlinkJedisPoolConfig.Builder().setHost(127.0.0.1).setPort(6379).build(); //数据写入redis dataStreamSource.addSink(new RedisSink(redisConf, new RedisMapperTuple2String, String() { Override public RedisCommandDescription getCommandDescription() { //指定redis命令,这里只演示最简单的设置字符串key return new RedisCommandDescription(RedisCommand.SET); } Override public String getKeyFromData(Tuple2String, String data) { //提取要存到redis的key return data.f0; } Override public String getValueFromData(Tuple2String, String data) { //提取要存到redis的value return data.f1; } })); env.execute(redis data sink); } }复制通过redis Cli 查看写入的数据127.0.0.1:6379 get A apple 127.0.0.1:6379 get B bird 127.0.0.1:6379 get C cat复制Redis 集群配置FlinkJedisClusterConfig config new FlinkJedisClusterConfig.Builder() .setNodes(new HashSetInetSocketAddress( Arrays.asList(new InetSocketAddress(host1, 6379), new InetSocketAddress(host2, 6379)))).build();复制Redis Sentinels配置FlinkJedisSentinelConfig sentinelConfig new FlinkJedisSentinelConfig.Builder() .setMasterName(master) .setSentinels(new HashSet(Arrays.asList(sentinel1, sentinel2))) .setPassword(12345) .setDatabase(1).build();复制4JDBCFlink官方提供了JDBC连接器只要引入连接器和mysql驱动即可。引入依赖dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-jdbc_2.11/artifactId version1.11.2/version /dependency dependency groupIdmysql/groupId artifactIdmysql-connector-java/artifactId version5.1.49/version /dependency复制flink应用代码public class JDBCSinkDemo { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); //准备数据源 ListProduct list new ArrayList(); list.add(new Product(跑鞋, 耐克, 998)); list.add(new Product(短裤, 李宁, 119)); list.add(new Product(袜子, 耐克, 68)); list.add(new Product(西服, 海澜之家, 1000)); DataStreamSourceProduct dataStreamSource env.fromCollection(list); String url jdbc:mysql://localhost:3306/flink_data?useUnicodetruecharacterEncodingutf-8serverTimezoneGMT; String sql insert into t_product(name, brand, price) values (?,?,?); dataStreamSource.addSink(JdbcSink.sink(sql, new JdbcStatementBuilderProduct() { Override public void accept(PreparedStatement ps, Product product) throws SQLException { ps.setString(1, product.getName()); ps.setString(2, product.getBrand()); ps.setDouble(3, product.getPrice()); } }, new JdbcConnectionOptions.JdbcConnectionOptionsBuilder() .withUrl(url) .withDriverName(com.mysql.jdbc.Driver) .withUsername(root) .withPassword(123456) .build())); env.execute(jdbc data sink); } }复制运行完毕查看数据库use flink_data; select * from t_product;复制5自定义除了Flink官方提供的第三方连接器外我们也可以自定义 Sink 来满足各种输出需求。自定义的 Sink需要直接或者间接实现 SinkFunction 接口一般直接继承抽象的富函数RichSinkFunction重写其open、close、invoke方法。相比于SinkFunction 富函数提供了操作生命周期的相关方法。需求实现一个自定义的sink将数据输出到Mysql数据库。MysqlSink.javaimport org.apache.flink.configuration.Configuration; import org.apache.flink.streaming.api.functions.sink.RichSinkFunction; import java.sql.Connection; import java.sql.DriverManager; import java.sql.PreparedStatement; public class MysqlSink extends RichSinkFunctionProduct { private PreparedStatement stmt; private Connection conn; Override public void open(Configuration parameters) throws Exception { super.open(parameters); String url jdbc:mysql://localhost:3306/flink_data?useUnicodetruecharacterEncodingutf-8serverTimezoneGMTautoReconnecttrue; String sql insert into t_product(name, brand, price) values (?,?,?); Class.forName(com.mysql.jdbc.Driver); conn DriverManager.getConnection(url, root, 123456); stmt conn.prepareStatement(sql); } Override public void close() throws Exception { super.close(); if (stmt ! null) { stmt.close(); } if (conn ! null) { conn.close(); } } Override public void invoke(Product product, Context context) throws Exception { stmt.setString(1, product.getName()); stmt.setString(2, product.getBrand()); stmt.setDouble(3, product.getPrice()); stmt.executeUpdate(); } }复制flink应用代码import org.apache.flink.streaming.api.datastream.DataStreamSource; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import java.util.ArrayList; import java.util.List; public class MysqlSinkDemo { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); //准备数据源 ListProduct list new ArrayList(); list.add(new Product(笔记本, 联想, 3998)); list.add(new Product(硬盘, 希捷, 219)); list.add(new Product(CPU, Intel, 668)); list.add(new Product(显示器, 飞利浦, 1400)); DataStreamSourceProduct dataStreamSource env.fromCollection(list); dataStreamSource.addSink(new MysqlSink()); env.execute(mysql data sink); } }