百度360必应搜狗淘宝本站头条
当前位置:网站首页 > 技术教程 > 正文

Flink教程-将流式数据写入redis

mhr18 2024-12-03 12:08 41 浏览 0 评论


  • 背景
  • 实例讲解
  • 动态生成key


背景

redis作为一个高吞吐的存储系统,在生产中有着广泛的应用,今天我们主要讲一下如何将流式数据写入redis,以及遇到的一些问题 解决。官方并没有提供写入redis的connector,所以我们采用apache的另一个项目bahir-flink [1]中提供的连接器来实现。

实例讲解

引入pom

 <dependency>
   <groupId>org.apache.flink</groupId>
   <artifactId>flink-connector-redis_2.11</artifactId>
   <version>1.1.5</version>
  </dependency>

构造数据源

这里我们主要是模拟一条用户信息

  //user,subject,province
  Tuple3<String,String,String> tuple = Tuple3.of("tom", "math", "beijing");
  DataStream<Tuple3<String,String,String>> dataStream = bsEnv.fromElements(tuple);

构造redis配置

  • 单机配置
 FlinkJedisPoolConfig conf = new FlinkJedisPoolConfig.Builder().setHost("10.160.85.185")
                                                                // 可选 .setPassword("1234")
                                                                .setPort(6379)
                                                                .build();
  • 集群配置
  InetSocketAddress host0 = new InetSocketAddress("host1", 6379);
  InetSocketAddress host1 = new InetSocketAddress("host2", 6379);
  InetSocketAddress host2 = new InetSocketAddress("host3", 6379);

  HashSet<InetSocketAddress> set = new HashSet<>();
  set.add(host0);
  set.add(host1);
  set.add(host2);

  FlinkJedisClusterConfig config = new FlinkJedisClusterConfig.Builder().setNodes(set)
                                                                        .build();

实现RedisMapper

我们需要实现一个RedisMapper接口的类,这个类的主要功能就是将我们自己的输入数据映射到redis的对应的类型。

我们看下RedisMapper接口,这里面总共有三个方法:

  • getCommandDescription:主要来获取我们写入哪种类型的数据,比如list、hash等等。
  • getKeyFromData:主要是从我们的输入数据中抽取key
  • getValueFromData:从我们的输入数据中抽取value
public interface RedisMapper<T> extends Function, Serializable {

 /**
  * Returns descriptor which defines data type.
  *
  * @return data type descriptor
  */
 RedisCommandDescription getCommandDescription();

 /**
  * Extracts key from data.
  *
  * @param data source data
  * @return key
  */
 String getKeyFromData(T data);

 /**
  * Extracts value from data.
  *
  * @param data source data
  * @return value
  */
 String getValueFromData(T data);
}

getCommandDescription方法返回一个RedisCommandDescription对象,我们看下RedisCommandDescription的构造方法:

 public RedisCommandDescription(RedisCommand redisCommand, String additionalKey) {
        ................
 }

 public RedisCommandDescription(RedisCommand redisCommand) {
  this(redisCommand, null);
 }

我们以数据写入hash结构为例,构造了一个key为HASH_NAME的RedisCommandDescription

new RedisCommandDescription(RedisCommand.HSET, "HASH_NAME");

两个构造方法区别就在于是否有第二个参数additionalKey,这个参数主要是针对SORTED_SET和HASH结构的,因为这两个结构需要有三个变量,其他的结构只需要两个变量就行了。

在hash结构里,这个additionalKey对应hash的key,getKeyFromData方法得到的数据对应hash的field,getValueFromData获取的数据对应hash的value。

最后我们数据写入对应的redis sink即可,写入的redis数据如下:

image

完整的代码请参考:

https://github.com/zhangjun0x01/bigdata-examples/blob/master/flink/src/main/java/connectors/redis/RedisSinkTest.java

动态生成key

我们看到,上面我们构造redis的hash结构的时候,key是写死的,也就是只能写入一个key,如果我的key是动态生成的,该怎么办呢?

比如我有一个类似的需求,流式数据是一些学生成绩信息,我的key想要的是学生的name,field是相应的科目,而value是这个科目对应的成绩。

目前系统没提供这样的功能,不过这个也没事,没有什么不是改源码解决不了的。

我们看下RedisSink中的invoke方法,

 public void invoke(IN input) throws Exception {
  String key = redisSinkMapper.getKeyFromData(input);
  String value = redisSinkMapper.getValueFromData(input);

  switch (redisCommand) {
      ....................
   case HSET:
    this.redisCommandsContainer.hset(this.additionalKey, key, value);
    break;
   default:
    throw new IllegalArgumentException("Cannot process such data type: " + redisCommand);
  }
 }

我们看到对于hash结构来说,key和value也就是从我们的RedisMapper的实现类中获取的,但是additionalKey却不是动态生成的,我们只需要改下这里。动态获取additionalKey就行。

public interface RedisMapper<T> extends Function, Serializable{

 RedisCommandDescription getCommandDescription();

 String getKeyFromData(T data);

 String getValueFromData(T data);

 String getAdditionalKey(T data);
}

我们给RedisMapper接口添加一个getAdditionalKey方法,然后在实现类中实现该方法。

然后在RedisSink的invoke方法动态获取additionalKey,修改源码之后的方法如下:

 @Override
 public void invoke(IN input) throws Exception {
  String key = redisSinkMapper.getKeyFromData(input);
  String value = redisSinkMapper.getValueFromData(input);
  String additionalKey = redisSinkMapper.getAdditionalKey(input);
  switch (redisCommand) {
         ..................
   case HSET:
    this.redisCommandsContainer.hset(additionalKey, key, value);
    break;
   default:
    throw new IllegalArgumentException("Cannot process such data type: " + redisCommand);
  }
 }

参考资料:
[1].https://github.com/apache/bahir-flink.git

更多内容,欢迎关注我的公众号【大数据技术与应用实战】

相关推荐

【推荐】一个开源免费、AI 驱动的智能数据管理系统,支持多数据库

如果您对源码&技术感兴趣,请点赞+收藏+转发+关注,大家的支持是我分享最大的动力!!!.前言在当今数据驱动的时代,高效、智能地管理数据已成为企业和个人不可或缺的能力。为了满足这一需求,我们推出了这款开...

Pure Storage推出统一数据管理云平台及新闪存阵列

PureStorage公司今日推出企业数据云(EnterpriseDataCloud),称其为组织在混合环境中存储、管理和使用数据方式的全面架构升级。该公司表示,EDC使组织能够在本地、云端和混...

对Java学习的10条建议(对java课程的建议)

不少Java的初学者一开始都是信心满满准备迎接挑战,但是经过一段时间的学习之后,多少都会碰到各种挫败,以下北风网就总结一些对于初学者非常有用的建议,希望能够给他们解决现实中的问题。Java编程的准备:...

SQLShift 重大更新:Oracle→PostgreSQL 存储过程转换功能上线!

官网:https://sqlshift.cn/6月,SQLShift迎来重大版本更新!作为国内首个支持Oracle->OceanBase存储过程智能转换的工具,SQLShift在过去一...

JDK21有没有什么稳定、简单又强势的特性?

佳未阿里云开发者2025年03月05日08:30浙江阿里妹导读这篇文章主要介绍了Java虚拟线程的发展及其在AJDK中的实现和优化。阅前声明:本文介绍的内容基于AJDK21.0.5[1]以及以上...

「松勤软件测试」网站总出现404 bug?总结8个原因,不信解决不了

在进行网站测试的时候,有没有碰到过网站崩溃,打不开,出现404错误等各种现象,如果你碰到了,那么恭喜你,你的网站出问题了,是什么原因导致网站出问题呢,根据松勤软件测试的总结如下:01数据库中的表空间不...

Java面试题及答案最全总结(2025版)

大家好,我是Java面试陪考员最近很多小伙伴在忙着找工作,给大家整理了一份非常全面的Java面试题及答案。涉及的内容非常全面,包含:Spring、MySQL、JVM、Redis、Linux、Sprin...

数据库日常运维工作内容(数据库日常运维 工作内容)

#数据库日常运维工作包括哪些内容?#数据库日常运维工作是一个涵盖多个层面的综合性任务,以下是详细的分类和内容说明:一、数据库运维核心工作监控与告警性能监控:实时监控CPU、内存、I/O、连接数、锁等待...

分布式之系统底层原理(上)(底层分布式技术)

作者:allanpan,腾讯IEG高级后台工程师导言分布式事务是分布式系统必不可少的组成部分,基本上只要实现一个分布式系统就逃不开对分布式事务的支持。本文从分布式事务这个概念切入,尝试对分布式事务...

oracle 死锁了怎么办?kill 进程 直接上干货

1、查看死锁是否存在selectusername,lockwait,status,machine,programfromv$sessionwheresidin(selectsession...

SpringBoot 各种分页查询方式详解(全网最全)

一、分页查询基础概念与原理1.1什么是分页查询分页查询是指将大量数据分割成多个小块(页)进行展示的技术,它是现代Web应用中必不可少的功能。想象一下你去图书馆找书,如果所有书都堆在一张桌子上,你很难...

《战场兄弟》全事件攻略 一般事件合同事件红装及隐藏职业攻略

《战场兄弟》全事件攻略,一般事件合同事件红装及隐藏职业攻略。《战场兄弟》事件奖励,事件条件。《战场兄弟》是OverhypeStudios制作发行的一款由xcom和桌游为灵感来源,以中世纪、低魔奇幻为...

LoadRunner(loadrunner录制不到脚本)

一、核心组件与工作流程LoadRunner性能测试工具-并发测试-正版软件下载-使用教程-价格-官方代理商的架构围绕三大核心组件构建,形成完整测试闭环:VirtualUserGenerator(...

Redis数据类型介绍(redis 数据类型)

介绍Redis支持五种数据类型:String(字符串),Hash(哈希),List(列表),Set(集合)及Zset(sortedset:有序集合)。1、字符串类型概述1.1、数据类型Redis支持...

RMAN备份监控及优化总结(rman备份原理)

今天主要介绍一下如何对RMAN备份监控及优化,这里就不讲rman备份的一些原理了,仅供参考。一、监控RMAN备份1、确定备份源与备份设备的最大速度从磁盘读的速度和磁带写的带度、备份的速度不可能超出这两...

取消回复欢迎 发表评论: