Flink-Connectors(连接器)(2)Redis
flink 提供了专门操作redis 的RedisSink,使用起来更方便,而且不用我们考虑性能的问题,接下来将主要介绍RedisSink 如何使用
https://bahir.apache.org/docs/flink/current/flink-streaming-redis/
必要依赖
<dependency>
<groupId>org.apache.bahir</groupId>
<artifactId>flink-connector-redis_2.11</artifactId>
<version>1.0</version>
<exclusions>
<exclusion>
<artifactId>flink-streaming-java_2.11</artifactId>
<groupId>org.apache.flink</groupId>
</exclusion>
<exclusion>
<artifactId>flink-runtime_2.11</artifactId>
<groupId>org.apache.flink</groupId>
</exclusion>
<exclusion>
<artifactId>flink-core</artifactId>
<groupId>org.apache.flink</groupId>
</exclusion>
<exclusion>
<artifactId>flink-java</artifactId>
<groupId>org.apache.flink</groupId>
</exclusion>
</exclusions>
</dependency>
RedisSink 的核心类是RedisMapper, 其是一个接口,使用时我们要编写自己的redis 操作类实现这个接口中的三个方法
1.getCommandDescription() :
设置使用的redis 数据结构类型,和key 的名称,通过RedisCommand 设置数据结构类型
2.String getKeyFromData(T data):
设置value 中的键值对key的值
3.String getValueFromData(T data);
设置value 中的键值对value的值
代码示例
package com.leilei;
import com.alibaba.fastjson.JSON;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import org.apache.flink.api.common.functions.FlatMapFunction;
import org.apache.flink.api.common.functions.MapFunction;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.common.typeinfo.Types;
import org.apache.flink.api.java.tuple.Tuple2;
import org.apache.flink.streaming.api.datastream.DataStreamSource;
import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.source.RichSourceFunction;
import org.apache.flink.streaming.connectors.redis.RedisSink;
import org.apache.flink.streaming.connectors.redis.common.config.FlinkJedisPoolConfig;
import org.apache.flink.streaming.connectors.redis.common.mapper.RedisCommand;
import org.apache.flink.streaming.connectors.redis.common.mapper.RedisCommandDescription;
import org.apache.flink.streaming.connectors.redis.common.mapper.RedisMapper;
import org.apache.flink.util.Collector;
/**
* @author lei
* @version 1.0
* @date 2021/3/14 14:24
* @desc flink redis-connector
*/
public class FlinkConnectorRedis {
public static void main(String[] args) {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStreamSource<VehicleAlarm> alarmSource = env.addSource(new AlarmSource());
SingleOutputStreamOperator<Tuple2<String, String>> source = alarmSource.map(new MapFunction<VehicleAlarm, VehicleAlarm>() {
@Override
public VehicleAlarm map(VehicleAlarm value) throws Exception {
value.setZone(value.getZone().toUpperCase());
return value;
}
}).flatMap(new FlatMapFunction<VehicleAlarm, Tuple2<String, String>>() {
@Override
public void flatMap(VehicleAlarm value, Collector<Tuple2<String, String>> out) throws Exception {
out.collect(Tuple2.of(value.getLicensePlate() + value.getZone(), JSON.toJSONString(value)));
}
});
// redis 连接器配置
FlinkJedisPoolConfig conf = new FlinkJedisPoolConfig.Builder()
.setHost("xxx")
.setDatabase(0)
.setPassword("lei123456").build();
// 输出到redis
source.addSink(new RedisSink<>(conf, new RedisExampleMapper()));
try {
env.execute("redis-sink-1");
} catch (Exception e) {
e.printStackTrace();
}
}
/**
* 输出到redis具体实现 (选择存入类型,以及指定KEY Value)
*/
public static class RedisExampleMapper implements RedisMapper<Tuple2<String, String>> {
/**
* redis 操作方式 RedisCommand.SET 则为普通String
* @return
*/
@Override
public RedisCommandDescription getCommandDescription() {
return new RedisCommandDescription(RedisCommand.SET);
}
/**
* 从数据中指定Key
* @param data
* @return
*/
@Override
public String getKeyFromData(Tuple2<String, String> data) {
return data.f0;
}
/**
* 从数据中指定value
* @param data
* @return
*/
@Override
public String getValueFromData(Tuple2<String, String> data) {
return data.f1;
}
}
@Data
@NoArgsConstructor
@AllArgsConstructor
public static class VehicleAlarm {
private String id;
private String licensePlate;
private String plateColor;
private Long deviceTime;
private String zone;
}
/**
* 自定义数据源
*/
public static class AlarmSource extends RichSourceFunction<VehicleAlarm> {
@Override
public void run(SourceContext<VehicleAlarm> ctx) throws Exception {
long id = System.currentTimeMillis() / 10000;
VehicleAlarm vehicleAlarm = new VehicleAlarm(String.valueOf(id), "川A" + id,
"紫", System.currentTimeMillis(), "sc");
ctx.collect(vehicleAlarm);
Thread.sleep(10000);
}
@Override
public void cancel() {
}
}
}