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); } }

相关新闻

免费AI视频增强神器Video2X:从模糊到高清的智能魔法

免费AI视频增强神器Video2X:从模糊到高清的智能魔法

免费AI视频增强神器Video2X:从模糊到高清的智能魔法 【免费下载链接】video2x A machine learning-based video super resolution and frame interpolation framework. Est. Hack the Valley II, 2018. 项目地址: https://gitcode.com/GitHub_Trending/vi/video2x…

2026/9/24 22:52:40 阅读更多 →
Python量化交易入门:从环境搭建到策略回测的完整实践指南

Python量化交易入门:从环境搭建到策略回测的完整实践指南

1. 这篇文章真正要解决的问题如果你对量化交易感兴趣,在B站、知乎、GitHub上搜索过相关教程,大概率会陷入一种困境:要么是过于理论化、堆砌数学公式的“劝退”视频,要么是只展示炫酷回测曲线、却不讲清楚如何从零开始的“魔术表演…

2026/9/23 22:59:33 阅读更多 →
高阶英语词汇测试系统设计与实现

高阶英语词汇测试系统设计与实现

1. 项目背景与核心价值 作为一名语言学习工具开发者,我注意到市面上大多数词汇测试App存在一个共同痛点:它们往往停留在基础词汇检测层面,缺乏对高阶学习者的针对性训练。这正是我们开发"超越英语八级词汇测试(101-300&#…

2026/9/24 0:06:50 阅读更多 →

最新新闻

工控现货:工业自动化备件的时效性与技术可靠性解析

工控现货:工业自动化备件的时效性与技术可靠性解析

1. “工控现货”不是电商标签,而是工业现场的生存语言“工控现货”这四个字,最近在自动化工程师的微信群、PLC维修论坛、甚至西门子/三菱授权服务商的报价单里出现频率陡增。它不是某个新出的电商平台栏目,也不是营销话术里的流量热词——它是…

2026/9/24 22:51:47 阅读更多 →
CodeBurn Menubar for Windows:基于 Tauri 2.x 的 AI 编码开销系统托盘应用开发指南

CodeBurn Menubar for Windows:基于 Tauri 2.x 的 AI 编码开销系统托盘应用开发指南

【免费下载链接】codeburn Free, local tool to track AI coding token usage and cost across 37 tools and agents (Claude Code, Cursor, Codex, Gemini and more), by model, project, and task. npx codeburn 项目地址: https://gitcode.com/gh_mirrors/co/cod…

2026/9/24 22:51:47 阅读更多 →
模糊人脸图像增强实战:UNetLike模型+感知损失+Django/Flask部署

模糊人脸图像增强实战:UNetLike模型+感知损失+Django/Flask部署

简介:本资源是一套高分通过的本科毕业设计项目,面向计算机、人工智能、软件工程等专业学生及初学者,聚焦模糊人脸图像增强这一典型CV任务,提供从模型训练到系统部署的完整实践方案。压缩包共23个文件,含9个核心Python源…

2026/9/24 22:51:47 阅读更多 →
iLeadE-588边缘计算盒子零代码可视化平台实战:从开箱到产线看板上线

iLeadE-588边缘计算盒子零代码可视化平台实战:从开箱到产线看板上线

工业现场做数据可视化,最怕两件事:一是为了看个实时曲线,得先装一堆客户端、配环境、调驱动,折腾半天还没看到数据;二是好不容易跑起来,想改个布局、加个图表,又得找开发排期。iLeadE-588 这台 …

2026/9/24 22:51:47 阅读更多 →
DeepAgents+MCP+A2A+Skills:多智能体系统从入门到生产实战

DeepAgents+MCP+A2A+Skills:多智能体系统从入门到生产实战

这几年做AI应用,我最大的感受是:单个模型再聪明,也扛不住真实业务的复杂度。真正把效率拉满的,反而是把任务拆开、让多个Agent各干各的,再用一套统一的协议把工具和协作串起来。DeepAgents MCP A2A Skills这套组合&…

2026/9/24 22:51:47 阅读更多 →
用curl验证vLLM与SGLang推理服务:从端口探活到接口响应全流程

用curl验证vLLM与SGLang推理服务:从端口探活到接口响应全流程

1. 先搞清楚:为什么服务"启动成功"不等于"能用"在 Ubuntu 上部署完 vLLM 或 SGLang 这类大模型推理引擎,日志里刷出Application startup complete或者Uvicorn running on http://0.0.0.0:8000的时候,很多人就觉得事情已经…

2026/9/24 22:50:46 阅读更多 →

日新闻

基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程

基于YOLOv8的渔船作业监控系统:从环境搭建到边缘部署全流程

简介:这是一套面向计算机、人工智能、自动化等专业学生与教师的毕业设计级项目资源,围绕YOLOv8实现渔船作业监控系统,可用于毕设、课程设计、大作业或项目立项演示。压缩包共97个文件,约24.21MB,以70个Python源码文件为…

2026/9/24 0:00:19 阅读更多 →
单细胞注释实战:基于Scanpy的标记基因与参考映射流程解析

单细胞注释实战:基于Scanpy的标记基因与参考映射流程解析

简介:一份基于单细胞RNA测序数据的细胞类型注释算法研究Python毕业设计源码,针对计算机相关专业正在做毕设或需要项目实战的学习者,可用于课程设计与期末大作业。项目代码完整、经导师指导评审通过,可直接运行,覆盖数据…

2026/9/24 0:00:19 阅读更多 →
C#源生成器实战:用增量生成器替代反射,告别AOT崩溃

C#源生成器实战:用增量生成器替代反射,告别AOT崩溃

第一次在项目里被反射卡住,是在一个老旧的WinForms模块里:几十个类依赖PropertyChanged通知,运行时反射读属性、发通知,每次启动慢半拍不说,一上.NET Native/AOT裁剪模式几乎全面崩盘。后来我把这段逻辑全部改成C#源生…

2026/9/24 0:00:19 阅读更多 →

周新闻

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

Flutter for OpenHarmony游戏卡片渐变背景实战:从原理到性能优化

直接铺开项目本身吧。这几个月我一直在折腾一件事:用Flutter给OpenHarmony做一款游戏集合类的App,说白了就是把若干小游戏塞进一个壳里,用统一入口分发。这个方向本身不算新鲜,真正让我花了不少心思的,是首页那堆游戏卡…

2026/9/24 14:34:13 阅读更多 →
Word表格编号全攻略:从列表编号到题注交叉引用

Word表格编号全攻略:从列表编号到题注交叉引用

写Word文档,最让人头疼的往往是那些“看起来不起眼”的小问题。比如表格编号这事:今天在表后面多加了两个空白行,明天给客户交稿前发现整个章节的编号全部错位,光是挨个改序号就能耗掉大半个下午。我前阵子帮人整理一份上百页的技…

2026/9/24 9:10:42 阅读更多 →
从第一个站到第二个站:独立开发者的静态网站选型与落地实践

从第一个站到第二个站:独立开发者的静态网站选型与落地实践

1. 项目概述1.1 核心需求解析做独立开发者这几年,说实话,第一个网站上线的那天晚上我兴奋得没睡着。但等它跑了半年,流量惨淡、功能臃肿、代码自己都懒得看第二遍之后,我才慢慢琢磨明白一个道理:第一个网站是练手&…

2026/9/24 14:33:56 阅读更多 →

月新闻

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能

持续集成 流水线自动化与 声明式交付 实践:原型怎样变成可用功能分类:[AI/大模型]细分主题:AI 增强型 CI/CD 流水线自动化与 GitOps 实践:Agent 工作流、工具调用与任务拆解:从原型到生产的验收清单很多团队在尝试用大…

2026/9/24 12:50:34 阅读更多 →
容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场

容器编排 生产环境运维与排障实战:复盘记录怎样真正派上用场分类:[工程技术]细分主题:Kubernetes 生产环境运维与排障实战:可复制的项目复盘模板与决策记录大部分团队的事故复盘报告,最后都变成了躺在 Confluence 或钉…

2026/9/24 14:33:48 阅读更多 →
容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步

容器 容器化技术与镜像安全管理:核心链路应该先拆哪一步分类:[工程技术]细分主题:Docker 容器化技术与镜像安全管理:核心链路的逐步实现与关键代码取舍面对一个积累了五六年历史包袱的单体架构应用(包含 Web 接口、后台…

2026/9/24 12:49:17 阅读更多 →