
大数据流处理后端【免费下载链接】stormApache Storm项目地址https://gitcode.com/gh_mirrors/storm6/storm点击查看免费下载本篇指南以 Apache Storm 仓库中external/storm-redis模块README为核心系统讲解如何在 Storm/Trident 拓扑中以 Redis 作为高速键值存储包括RedisStoreBolt/RedisLookupBolt/RedisFilterBolt三个基础 Bolt、自定义 Bolt 基类AbstractRedisBolt、RedisDataTypeDescription数据类型映射以及面向 Trident 的RedisState/RedisClusterState两套状态实现。读完本文你将能独立完成“单词计数落库 Redis”“按 Redis 黑名单过滤元组”“Redis 键值实时查询并外发”等典型拓扑并理解每一层在源码中的真实执行逻辑。模块概览storm-redis 是什么storm-redis 是 Apache Storm 的官方外部模块为 Storm 与 Trident 提供基于 Redis 的读写集成底层客户端使用 Jedis。它位于仓库的 external/storm-redis 目录模块坐标定义在 pom.xml 中依赖storm-client、jedis、guava以及 Jackson用于 JSON 序列化场景测试阶段还引入 testcontainers 以便在真实 Redis 容器上验证行为。在 Maven 项目中引入方式如下版本号与当前仓库一致使用${storm.version}占位符替换为实际 Storm 版本dependency groupIdorg.apache.storm/groupId artifactIdstorm-redis/artifactId version${storm.version}/version typejar/type /dependency从源码结构看模块被清晰地划分为四层见 src/main/java/org/apache/storm/redis包路径职责bolt基础 Bolt 与可继承的抽象 Boltcommon.configJedisPoolConfig单机与JedisClusterConfig集群连接配置common.mapper将 Tuple 映射为 Redis key/value 的 Mapper 接口族common.container统一封装Jedis/JedisCluster的连接容器trident.stateTrident 的 State / StateUpdater / StateQuerier 实现核心抽象Mapper 数据类型描述storm-redis 的顶层设计思想是**“一个 Tuple 对应一对 key/value匹配规则由 TupleMapper 定义”**。所有 Mapper 都继承自RedisMapperRedisMapper.java核心方法是返回数据类型描述public interface RedisMapper { RedisDataTypeDescription getDataTypeDescription(); }RedisDataTypeDescription 与支持的数据类型RedisDataTypeDescriptionRedisDataTypeDescription.java由两部分组成dataType数据类型与additionalKey附加键。其枚举定义了七种 Redis 数据类型public enum RedisDataType { STRING, HASH, LIST, SET, SORTED_SET, HYPER_LOG_LOG, GEO }重要约束源码中强制校验当 dataType 为HASH、SORTED_SET或GEO时构造函数会在additionalKey null时抛出IllegalArgumentException(Hash, Sorted Set and GEO should have additional key)。因为这些类型中“从 Tuple 转换出的键”只是元素field/member真正的 Redis 主键必须由 additionalKey 提供。而SET在RedisFilterBolt场景下同样要求提供 additionalKey。Mapper 接口族三个基础 Bolt 分别对应三种 Mapper均扩展了RedisMapperRedisStoreMapper写getKeyFromTuplegetValueFromTuple从 Tuple 提取待写入的 key 与 valueRedisLookupMapper查getKeyFromTuple提取查询键toTuple(ITuple input, Object value)把 Redis 返回值组装成输出 TupleRedisFilterMapper过滤getKeyFromTuple提取待判断的键declareOutputFields()声明的字段必须与输入流一致因为过滤器命中时会把原始输入 Tuple 原样转发。基础 Bolt 之一RedisStoreBolt写入 RedisRedisStoreBoltRedisStoreBolt.java对每个输入 Tuple 提取 key/value 后按数据类型执行对应 Jedis 命令成功后ack(input)异常则reportErrorfail(input)。其内部switch展示了各数据类型真实落库命令dataType执行的 Redis 命令说明STRINGset(key, value)普通字符串键值LISTrpush(key, value)追加到列表尾部HASHhset(additionalKey, key, value)写入 hash 的 fieldSETsadd(key, value)写入集合成员SORTED_SETzadd(additionalKey, Double.valueOf(value), key)value 需为可解析的 double 分数HYPER_LOG_LOGpfadd(key, value)基数统计GEOgeoadd(additionalKey, longitude, latitude, key)value 必须形如longitude:latitude否则抛异常完整示例WordCountStoreMapperclass WordCountStoreMapper implements RedisStoreMapper { private RedisDataTypeDescription description; private final String hashKey wordCount; public WordCountStoreMapper() { description new RedisDataTypeDescription( RedisDataTypeDescription.RedisDataType.HASH, hashKey); } Override public RedisDataTypeDescription getDataTypeDescription() { return description; } Override public String getKeyFromTuple(ITuple tuple) { return tuple.getStringByField(word); } Override public String getValueFromTuple(ITuple tuple) { return tuple.getStringByField(count); } }接入拓扑JedisPoolConfig poolConfig new JedisPoolConfig.Builder() .setHost(host).setPort(port).build(); RedisStoreMapper storeMapper new WordCountStoreMapper(); RedisStoreBolt storeBolt new RedisStoreBolt(poolConfig, storeMapper);最终效果单词word作为 field、计数count作为 value 写入 Redis hashwordCount中即HSET wordCount word count。基础 Bolt 之二RedisLookupBolt查询 RedisRedisLookupBoltRedisLookupBolt.java用getKeyFromTuple得到查询键按数据类型调用命令再把结果交给lookupMapper.toTuple(input, lookupValue)转成输出 Tuple 发出支持一次返回多个 Values。各数据类型的查询命令同样可从源码的switch确认dataType执行的 Redis 命令返回语义STRINGget(key)键值LISTlpop(key)弹出列表头部元素HASHhget(additionalKey, key)hash 中指定 field 的值SETscard(key)集合基数SORTED_SETzscore(additionalKey, key)成员分数HYPER_LOG_LOGpfcount(key)基数估算GEOgeopos(additionalKey, key)经纬度坐标完整示例WordCountRedisLookupMapperclass WordCountRedisLookupMapper implements RedisLookupMapper { private RedisDataTypeDescription description; private final String hashKey wordCount; public WordCountRedisLookupMapper() { description new RedisDataTypeDescription( RedisDataTypeDescription.RedisDataType.HASH, hashKey); } Override public ListValues toTuple(ITuple input, Object value) { String member getKeyFromTuple(input); ListValues values Lists.newArrayList(); values.add(new Values(member, value)); return values; } Override public void declareOutputFields(OutputFieldsDeclarer declarer) { declarer.declare(new Fields(wordName, count)); } Override public RedisDataTypeDescription getDataTypeDescription() { return description; } Override public String getKeyFromTuple(ITuple tuple) { return tuple.getStringByField(word); } Override public String getValueFromTuple(ITuple tuple) { return null; } }接入拓扑JedisPoolConfig poolConfig new JedisPoolConfig.Builder() .setHost(host).setPort(port).build(); RedisLookupMapper lookupMapper new WordCountRedisLookupMapper(); RedisLookupBolt lookupBolt new RedisLookupBolt(poolConfig, lookupMapper);基础 Bolt 之三RedisFilterBolt过滤元组RedisFilterBoltRedisFilterBolt.java的作用是当键或字段在 Redis 中不存在时过滤掉该 Tuple存在时把输入 Tuple 原样转发到默认流。其判断逻辑按数据类型而异STRINGexists(key)检查键空间是否存在HASH/SORTED_SET/GEO检查字段/成员是否存在于该数据结构hexists/zrank非空 /geopos非空SET/HYPER_LOG_LOG检查值是否存在于该数据结构sismember/pfcount 0。两个源码级注意点构造函数中若 dataType 为SET而 additionalKey 为 null会直接抛出IllegalArgumentException(additionalKey should be defined)若只想“不论数据类型、单纯判断键是否存在”可将 dataType 指定为STRING此时底层走EXISTS命令与键的真实类型无关。完整示例黑名单过滤class BlacklistWordFilterMapper implements RedisFilterMapper { private RedisDataTypeDescription description; private final String setKey blacklist; public BlacklistWordFilterMapper() { description new RedisDataTypeDescription( RedisDataTypeDescription.RedisDataType.SET, setKey); } Override public void declareOutputFields(OutputFieldsDeclarer declarer) { declarer.declare(new Fields(word, count)); } Override public RedisDataTypeDescription getDataTypeDescription() { return description; } Override public String getKeyFromTuple(ITuple tuple) { return tuple.getStringByField(word); } Override public String getValueFromTuple(ITuple tuple) { return null; } }接入拓扑JedisPoolConfig poolConfig new JedisPoolConfig.Builder() .setHost(host).setPort(port).build(); RedisFilterMapper filterMapper new BlacklistWordFilterMapper(); RedisFilterBolt filterBolt new RedisFilterBolt(poolConfig, filterMapper);该功能的测试实现于 RedisFilterBoltTest.java测试通过 testcontainers 拉起redis:7.2.3-alpine真实容器对 STRING/HASH/SET/SORTED_SET/HYPER_LOG_LOG/GEO 各数据类型逐一验证“键不存在则元组被过滤、键存在则元组被转发”可作为理解RedisFilterBolt行为的直接参考。连接配置单机 JedisPoolConfig 与集群 JedisClusterConfig三个基础 Bolt 均提供两套构造器传入JedisPoolConfig走单机 JedisPool传入JedisClusterConfig走 Redis Cluster。JedisPoolConfig单机定义于 JedisPoolConfig.javaBuilder 默认值取自 Jedis 的Protocol常量Builder 方法含义默认值setHost(String)主机名或 IPProtocol.DEFAULT_HOSTsetPort(int)端口Protocol.DEFAULT_PORT6379setTimeout(int)socket/连接超时毫秒Protocol.DEFAULT_TIMEOUTsetDatabase(int)数据库索引Protocol.DEFAULT_DATABASE0setPassword(String)密码可选无JedisClusterConfigRedis Cluster定义于 JedisClusterConfig.javaBuilder 方法含义默认值setNodes(SetInetSocketAddress)节点列表必填无缺失抛NullPointerExceptionsetTimeout(int)socket/连接超时Protocol.DEFAULT_TIMEOUTsetMaxRedirections(int)跟随 MOVED/ASK 重定向的次数上限5setPassword(String)密码可选无自定义 Bolt继承 AbstractRedisBolt当RedisStoreBolt/RedisLookupBolt/RedisFilterBolt无法覆盖业务场景时可直接继承AbstractRedisBoltAbstractRedisBolt.java。它基于BaseTickTupleAwareRichBolt在prepare()阶段根据传入的JedisPoolConfig或JedisClusterConfig构建连接容器二者皆未设置则抛IllegalArgumentException(Jedis configuration not found)并在cleanup()时关闭容器。推荐编码模式源码 Javadoc 明确给出借助getInstance()借用 Jedis 命令句柄业务处理完后在finally中调用returnInstance()归还并自行负责 ack/failJedisCommands jedisCommands null; try { jedisCommands getInstance(); // do some works } finally { if (jedisCommands ! null) { returnInstance(jedisCommands); } }完整示例自定义查询总词频的 Boltpublic static class LookupWordTotalCountBolt extends AbstractRedisBolt { private static final Logger LOG LoggerFactory.getLogger(LookupWordTotalCountBolt.class); private static final Random RANDOM new Random(); public LookupWordTotalCountBolt(JedisPoolConfig config) { super(config); } public LookupWordTotalCountBolt(JedisClusterConfig config) { super(config); } Override public void execute(Tuple input) { JedisCommands jedisCommands null; try { jedisCommands getInstance(); String wordName input.getStringByField(word); String countStr jedisCommands.get(wordName); if (countStr ! null) { int count Integer.parseInt(countStr); this.collector.emit(new Values(wordName, count)); // print lookup result with low probability if(RANDOM.nextInt(1000) 995) { LOG.info(Lookup result - word : wordName / count : count); } } else { // skip LOG.warn(Word not found in Redis - word : wordName); } } finally { if (jedisCommands ! null) { returnInstance(jedisCommands); } this.collector.ack(input); } } Override public void declareOutputFields(OutputFieldsDeclarer declarer) { // wordName, count declarer.declare(new Fields(wordName, count)); } }注意此处getInstance()返回的JedisCommandsContainer对外暴露的是单键操作能力源码中的RedisCommands接口RedisCommands.java统一了Jedis与JedisCluster的二进制/字符串方法差异这正是“同一定义、单机与集群两套实现”的底层依据。Trident State 集成Trident 场景下storm-redis 提供两套 State 体系RedisState/RedisMapState基于Jedis接口面向单机 RedisRedisClusterState/RedisClusterMapState基于JedisCluster接口面向 Redis Cluster。单机RedisStateRedisState.FactoryRedisState.java在makeState()中依据JedisPoolConfig创建JedisPool并包装成 StateRedisStateUpdaterRedisStateUpdater.java使用Pipeline批量写入一批(key, value)对RedisStateQuerierRedisStateQuerier.java使用mget/hmget批量查询。JedisPoolConfig poolConfig new JedisPoolConfig.Builder() .setHost(redisHost).setPort(redisPort) .build(); RedisStoreMapper storeMapper new WordCountStoreMapper(); RedisLookupMapper lookupMapper new WordCountLookupMapper(); RedisState.Factory factory new RedisState.Factory(poolConfig); TridentTopology topology new TridentTopology(); Stream stream topology.newStream(spout1, spout); stream.partitionPersist(factory, fields, new RedisStateUpdater(storeMapper).withExpire(86400000), new Fields()); TridentState state topology.newStaticState(factory); stream stream.stateQuery(state, new Fields(word), new RedisStateQuerier(lookupMapper), new Fields(columnName,columnValue));集群RedisClusterStateSetInetSocketAddress nodes new HashSetInetSocketAddress(); for (String hostPort : redisHostPort.split(,)) { String[] host_port hostPort.split(:); nodes.add(new InetSocketAddress(host_port[0], Integer.valueOf(host_port[1]))); } JedisClusterConfig clusterConfig new JedisClusterConfig.Builder().setNodes(nodes) .build(); RedisStoreMapper storeMapper new WordCountStoreMapper(); RedisLookupMapper lookupMapper new WordCountLookupMapper(); RedisClusterState.Factory factory new RedisClusterState.Factory(clusterConfig); TridentTopology topology new TridentTopology(); Stream stream topology.newStream(spout1, spout); stream.partitionPersist(factory, fields, new RedisClusterStateUpdater(storeMapper).withExpire(86400000), new Fields()); TridentState state topology.newStaticState(factory); stream stream.stateQuery(state, new Fields(word), new RedisClusterStateQuerier(lookupMapper), new Fields(columnName,columnValue));withExpire 的过期语义源码细节withExpire(expireIntervalSec)与setExpireInterval(int)见 AbstractRedisStateUpdater.java仅在传入值大于 0 时生效。底层行为RedisStateUpdater.java分为两类STRING逐 key 执行setex(key, expireIntervalSec, value)HASH批量hset后只对 additionalKey 整体执行一次expire——即过期的是整个 hash 键本身而非单个 field源码注释明确提醒“expires key itself entirely, so use it with caution”会让整个键彻底过期使用时务必谨慎。端到端示例PersistentWordCount仓库的示例工程 examples/storm-redis-examples 中提供了可直接参考的完整拓扑 PersistentWordCount.java数据流为WordSpout → WordCounter → RedisStoreBoltfieldsGrouping保证同词计数、shuffleGrouping分发到存储 Bolt最终把(word, count)写入 Redis hash// wordSpout countBolt RedisBolt TopologyBuilder builder new TopologyBuilder(); builder.setSpout(WORD_SPOUT, spout, 1); builder.setBolt(COUNT_BOLT, bolt, 1).fieldsGrouping(WORD_SPOUT, new Fields(word)); builder.setBolt(STORE_BOLT, storeBolt, 1).shuffleGrouping(COUNT_BOLT); Config config new Config(); StormSubmitter.submitTopology(topoName, config, builder.createTopology());提交命令支持参数化 Redis 地址与拓扑名Usage: PersistentWordCount redis host redis port (topology name)同目录下的 LookupWordCount.java查询已入库计数、WhitelistWordCount.java白名单过滤以及 Trident 四件套 WordCountTridentRedis.java、WordCountTridentRedisCluster.java、WordCountTridentRedisMap.java、WordCountTridentRedisClusterMap.java共同构成了从“普通 Bolt”到“Trident State”的完整实践闭环建议对照源码逐一研读。使用前提与限制本模块依赖 Jedis需保证运行时 classpath 包含jedis及storm-client见 pom.xml连接配置对象JedisPoolConfig/JedisClusterConfig均实现Serializable可随拓扑分发到 Worker集群配置的节点集合不可为空HASH/SORTED_SET/GEO三种数据类型强制要求additionalKeyRedisFilterBolt使用SET类型时同样强制要求否则构造即抛异常GEO写入时 value 必须为longitude:latitude格式以冒号分隔、两段SORTED_SET的 value 必须可解析为 double 分数Trident State 的 HASH 过期作用于整个键属于“整体过期”语义需结合业务谨慎使用。Apache Storm 与 Redis 的组合让实时计算链路的“计算”与“存储/过滤”在同一个数据流内无缝衔接——无论是计数落库、黑白名单实时过滤还是 Trident 有状态批处理storm-redis都以一套统一的数据类型映射抽象屏蔽了单机与集群的差异值得作为集成外部存储的参考范式。赞分享大数据流处理后端【免费下载链接】stormApache Storm项目地址https://gitcode.com/gh_mirrors/storm6/storm点击查看免费下载相关推荐Apache Pulsar 与 Apache Storm 集成Pulsar Storm Adaptor 的 Spout/Bolt 实战指南Apache Pulsar 与 Apache Storm 集成Pulsar Storm Adaptor 的 Spout/Bolt 实战指南 本指南以 Puls消息队列后端流处理Apache Pulsar 与 Apache Storm 集成指南Pulsar Storm Adaptor 的 Spout 与 Bolt 实战Apache Pulsar 与 Apache Storm 集成指南Pulsar Storm Adaptor 的 Spout 与 Bolt 实战 导读 本文讲解消息队列后端流处理Apache Pulsar 与 Apache Storm 集成指南Pulsar Storm Adaptor 的 Spout 与 Bolt 完整实战Apache Pulsar 与 Apache Storm 集成指南Pulsar Storm Adaptor 的 Spout 与 Bolt 完整实战 本篇技术指消息队列后端流处理上一篇5分钟学会智慧树自动刷课终极Chrome插件安装与使用指南下一篇在 Entity Framework Core 中使用 SQL Server JSON 字段Blogging 示例实战解析创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考