ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Storm External 模块生态:连接器架构、治理机制与 15 个集成模块全景指南

Apache Storm External 模块生态:连接器架构、治理机制与 15 个集成模块全景指南 大数据流处理后端【免费下载链接】stormApache Storm项目地址https://gitcode.com/gh_mirrors/storm6/storm点击查看免费下载Apache Storm 的核心运行时由storm-client、storm-core、storm-server等模块构成而external/目录下则聚集了一批官方扩展件——它们不是 Storm 运行所必需的却能让拓扑便捷地对接 Kafka、HDFS、JDBC、Redis、JMS、Iceberg 等大数据生态中高频使用的技术。本文以 external/README.md 为骨架结合仓库中 15 个外部模块的真实 README、pom.xml与示例工程系统讲解 external 模块的定位、版本兼容策略、Committer Sponsor 治理机制并逐一给出各模块的能力清单、关键 API 与可落地的配置/命令示例帮助读者在引入连接器时快速判断选型、正确打包并规避常见的 classpath 与版本坑。什么是 externalStorm 的官方扩展模块体系在 Apache Storm 仓库中external/是一个承载扩展模块的聚合目录。根据 external/README.md 的定义external 是一组不为 Storm 运行所必需、但能扩展 Storm 能力的模块其主要价值在于提供与其他常用技术的集成功能。也就是说核心运行时不依赖它们移除或缺失任何 external 模块Nimbus、Supervisor、Worker 与拓扑的基本执行不受影响它们承载连接器职责把 Storm 与消息队列Kafka、JMS、数据库JDBC、Redis、文件系统HDFS、数据湖Iceberg等系统打通它们是官方维护的由 Apache Storm 项目组统一发布而不是第三方松散组件。从构建产物看external/pom.xml 是一个packagingpom的聚合模块artifactIdstorm-external/artifactId描述为 Aggregator for Storm external integrations其中显式声明了 15 个子模块modules modulestorm-autocreds/module modulestorm-blobstore-migration/module modulestorm-hdfs/module modulestorm-hdfs-blobstore/module modulestorm-hdfs-oci/module modulestorm-iceberg/module modulestorm-jdbc/module modulestorm-jms/module modulestorm-kafka-client/module modulestorm-kafka-migration/module modulestorm-kafka-monitor/module modulestorm-metrics/module modulestorm-metrics-prometheus/module modulestorm-redis/module /modules这套模块划分把流计算内核与外部集成在工程上做了清晰解耦内核迭代不影响连接器连接器演进也不污染内核二者通过同一 Maven 坐标org.apache.storm:*与版本号保持一致。版本兼容策略与 Storm 同步发布external 模块的一个重要约束是与 Storm 主版本同步发布released in tandem with Storm目的是维持版本兼容性。这意味着连接器与内核共享同一版本号当前仓库版本为3.1.1-SNAPSHOT见 external/pom.xml 的parent定义当你引入某个连接器时应使用与 Storm 完全一致的版本避免因 ABI 或内部 API 变化导致的运行时冲突升级 Storm 时external 模块通常也需同步升级。这一策略在配套示例工程中体现得最直观。例如 examples/storm-kafka-client-examples/pom.xml、examples/storm-jdbc-examples/pom.xml、examples/storm-redis-examples/pom.xml 等都通过${storm.version}属性引用连接器版本与主工程保持同频。实际项目中推荐同样以属性占位的方式声明依赖例如dependency groupIdorg.apache.storm/groupId artifactIdstorm-kafka-client/artifactId version${storm.version}/version /dependencyCommitter Sponsor外部模块的防代码腐烂治理机制原文档的另一个核心概念是Committer Sponsor提交者赞助人。所谓 Committer Sponsor就是一位对某个模块表示过支持意愿的 Apache Storm Committer。项目组希望每个模块至少有一位 Sponsor以此在一定程度上防范 code rot代码腐烂和 abandonware弃维护风险——即模块长期无人修复、无人跟进而逐渐失效。需要特别澄清的是Committer Sponsor并不拥有任何特殊角色、特权或义务。Apache Storm Committers 对整个代码库拥有平等的权力和责任Sponsor 本质上只是表达了我对这个模块感兴趣愿意在力所能及的地方帮忙的意愿。这一点从仓库中多个模块 README 末尾的署名可以印证例如 external/storm-kafka-monitor/README.md 与 external/storm-jms/README.markdown 均列出 Committer Sponsors 一节。对使用者而言这意味着每个 external 模块都有人照看遇到问题可以通过社区协作推进修复但模块的维护强度、活跃度并不因 Sponsor 的存在而得到硬性承诺选型时仍应结合自身场景评估。数据接入类连接器把消息与键值数据引入拓扑storm-kafka-client基于新 consumer API 的 Kafka 集成external/storm-kafka-client/README.md 明确指出该模块包含使用新 Apache Kafka consumer API 的 Spout 与 Bolt用于通过 kafka-client 库读写 Kafka其完整使用说明见 docs/storm-kafka-client.md。这是当前推荐的 Kafka 接入方式适用于消费/生产双场景仓库在 examples/storm-kafka-client-examples 提供了可直接参考的 Spout/Bolt 用法与测试。配套的 external/storm-kafka-migration/README.md 则解决了从旧版storm-kafka迁移时的偏移量衔接问题迁移普通非 Trident偏移量运行java -cp * org.apache.storm.kafka.migration.KafkaSpoutMigration your-config-file.yaml工具会把偏移量迁入 Kafka 指定 consumer group迁移 Trident 偏移量运行java -cp * org.apache.storm.kafka.migration.KafkaTridentSpoutMigration your-config-file.yaml偏移量写入配置指定的 Zookeeper 路径并要求TridentDataSource.newStream的 txid 与配置中的new.topology.txid保持一致。迁移前需要停止拓扑且 jar 需与匹配 broker 版本的org.apache.kafka:kafka-clients放在同一目录。示例配置见该模块的 src/main/conf 目录。storm-jms数据无关的 JMS 桥接框架external/storm-jms/README.markdown 将 storm-jms 定位为在 Storm 框架内集成 JMS 消息传递的通用框架通过通用 JMS Spout 把 JMS 消息注入 Storm通过通用 JMS Bolt 把拓扑数据发布到 JMS 目的地topic 或 queue。两个组件都保持数据无关data agnostic——领域逻辑由用户提供的桥接类封装具体用法与示例参见 docs/storm-jms.md 及 examples/storm-jms-examples。storm-redis基于 Jedis 的 Redis 读写/过滤external/storm-redis/README.md 说明该模块基于 Jedis 客户端提供三类基础 BoltRedisLookupBolt按 key 从 Redis 取值RedisStoreBolt把 key/value 写入 RedisRedisFilterBolt过滤掉 key 或 field 在 Redis 上不存在的 tuple。使用上通过TupleMapper定义 tuple 到 key/value 的匹配规则通过RedisDataTypeDescription选择数据类型HASH、SET、SORTED SET 等部分类型需要额外的 keytuple 转换出的 key 成为元素并分别搭配RedisLookupMapper、RedisStoreMapper、RedisFilterMapper使用。其中 FilterBolt 会转发输入 tuple因此实现RedisFilterMapper时declareOutputFields()需声明与输入流相同的字段。一个典型用法是来自模块 README 的 WordCount 示例骨架class 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)); } // ... }完整文档见 docs/storm-redis.md示例工程见 examples/storm-redis-examples。数据输出/存储类连接器把流写入文件、数据库与数据湖storm-hdfs可定制化极强的 HDFS 写入组件external/storm-hdfs/README.md 提供了 HDFS Bolt 与 HDFS Spout。以 Bolt 为例一个管道分隔、每 1000 条 tuple 同步一次、文件到 5MB 轮转的最小配置如下// 字段分隔符用 | 而不是 , RecordFormat format new DelimitedRecordFormat() .withFieldDelimiter(|); // 每 1k 条 tuple 同步一次文件系统 SyncPolicy syncPolicy new CountSyncPolicy(1000); // 文件达到 5MB 时轮转 FileRotationPolicy rotationPolicy new FileSizeRotationPolicy(5.0f, Units.MB); FileNameFormat fileNameFormat new DefaultFileNameFormat() .withPath(/foo/); HdfsBolt bolt new HdfsBolt() .withFsUrl(hdfs://localhost:54310) .withFileNameFormat(fileNameFormat) .withRecordFormat(format) .withRotationPolicy(rotationPolicy) .withSyncPolicy(syncPolicy);该模块的可扩展点通过一组Serializable接口暴露使用者可自由定制RecordFormatbyte[] format(Tuple tuple)控制行格式内置DelimitedRecordFormat可生成 CSV、制表符分隔等格式FileNameFormatprepare/getName/getPath控制命名内置DefaultFileNameFormat生成{prefix}{componentId}-{taskId}-{rotationNum}-{timestamp}{extension}如MyBolt-5-7-1390579837830.txt默认 prefix 为空、扩展名为.txt更新的SimpleFileNameFormat支持$TIME、$NUM、$HOST、$COMPONENT、$TASK等占位符Trident 版另有$PARTITION例如seq.$TIME.$HOST.$COMPONENT.$NUM.dat默认名为$TIME.$NUM.txt默认时间格式yyyyMMddHHmmssSyncPolicymark/reset决定缓冲数据何时 flush 到文件系统CountSyncPolicy按 tuple 数触发FileRotationPolicymark/reset决定数据文件何时轮转FileSizeRotationPolicy按大小触发。README 还强调了一个高频踩坑点——打包必须用 maven-shade-plugin 而非 maven-assembly-pluginshade 插件能合并 JAR manifest 条目供 Hadoop client 做 URL scheme 解析。若出现java.lang.RuntimeException: Error preparing HdfsBolt: No FileSystem for scheme: hdfs基本可以断定拓扑 jar 打包方式不正确。推荐的 shade 配置需带上ServicesResourceTransformer与ManifestResourceTransformer。另外默认 Hadoop 依赖为hadoop-client与hadoop-hdfs2.6.1排除slf4j-log4j12如使用其他 Hadoop 版本应在自身 pom 中排除并替换版本不兼容常表现为com.google.protobuf.InvalidProtocolBufferException: Protocol message contained an invalid tag (zero)。完整文档见 docs/storm-hdfs.md示例见 examples/storm-hdfs-examples。storm-jdbc面向单表的 JDBC 写入与查询external/storm-jdbc/README.md 覆盖 Storm/Trident 的 JDBC 集成既可以把 tuple 插入数据库表也可以执行 select 查询回填enrichtuple。核心抽象有三层ConnectionProviderprepare/getConnection/cleanup要求幂等连接池抽象开箱即用HikariCPConnectionProviderJdbcMapperListColumn getColumns(ITuple tuple)定义 tuple 到数据库行的映射返回列表的顺序必须与 SQL 中占位符顺序一致——连接器不解析 SQL 按列名匹配占位符因此也能适配 Phoenix 这类仅支持 upsert 的非标准 SQL 框架JdbcInsertBolt通过withTableName或withInsertQuery指定目标可选withQueryTimeoutSecs设置查询超时默认取topology.message.timeout.secs-1表示不设超时建议设为不超过消息超时值。Map hikariConfigMap Maps.newHashMap(); hikariConfigMap.put(dataSourceClassName,com.mysql.jdbc.jdbc2.optional.MysqlDataSource); hikariConfigMap.put(dataSource.url, jdbc:mysql://localhost/test); hikariConfigMap.put(dataSource.user,root); hikariConfigMap.put(dataSource.password,password); ConnectionProvider connectionProvider new HikariCPConnectionProvider(hikariConfigMap); String tableName user_details; JdbcMapper simpleJdbcMapper new SimpleJdbcMapper(tableName, connectionProvider); JdbcInsertBolt userPersistenceBolt new JdbcInsertBolt(connectionProvider, simpleJdbcMapper) .withTableName(user) .withQueryTimeoutSecs(30); // 或.withInsertQuery(insert into user values (?,?))更多用法见 docs/storm-jdbc.md 与 examples/storm-jdbc-examples。storm-iceberg直达数据湖的原子提交写入external/storm-iceberg/README.md 描述了一个直接从拓扑写入 Apache Iceberg 表的 Bolt中间不再需要 Kafka Connect 或 Spark 作业具备原子提交与至少一次at-least-once投递语义读取方永远不会看到半截批次tuple 只有在包含它的 commit 落地后才会被 ack因此不会丢失被重放的批次会被重写由于 sink 是 append-only 且不写 equality delete重复行会保留直到下游清理。目标表必须为 append-only 且为 format version 2。典型用法MapString, String catalogProps new HashMap(); catalogProps.put(type, rest); catalogProps.put(uri, http://rest-catalog:8181); IcebergOptions options new IcebergOptions.Builder() .withCatalogProperties(catalogProps) .withTable(db.events) .withCommitIntervalBytes(128L * 1024 * 1024) .build(); TopologyBuilder builder new TopologyBuilder(); builder.setSpout(events, spout, 2); builder.setBolt(iceberg, new IcebergBolt(options), 4);完整说明见 docs/storm-iceberg.md示例见 examples/storm-iceberg-examples。安全与凭据storm-autocreds 自动获取、分发与续期令牌external/storm-autocreds/README.md 描述了面向安全集群的凭据自动化能力让 Storm自动获取、分发并续期 Hadoop delegation token使拓扑无需在每个 worker 主机上分发 keytab 即可访问开启 Kerberos 的 HDFS/HBase。工作流为拓扑提交时Nimbus代表提交用户获取 delegation token 并随拓扑下发Worker将 token 解包到自身Subject/UserGroupInformationNimbus为长跑拓扑周期性续期 token。由于这些插件运行在守护进程 classpathNimbus/Supervisor上且会拉入完整的 Hadoop/HBase 客户端依赖树只有安全 Hadoop 部署才需要因此完整版发行包apache-storm-x.x.x.tar.gz打包了 jarlite 发行包则不打包。安装到$STORM_HOME/extlib-daemon后需重启 Nimbus 与 Supervisors。lite 发行包可用辅助脚本从 Maven Central 解析完整依赖闭包# 自动识别版本安装 $STORM_HOME/bin/storm-autocreds-fetch # 指定版本与目标目录 bin/storm-autocreds-fetch --version 3.0.0 --dest /opt/storm/extlib-daemon # 透传 Maven 参数内网镜像 / 离线仓库 bin/storm-autocreds-fetch -- -s /etc/maven/settings.xml bin/storm-autocreds-fetch -- -Dmaven.repo.local/srv/offline-repo -o对应的storm.yaml配置分两侧带Nimbus后缀的类在 Nimbus 侧获取与续期 token不带后缀的类在 worker 侧解包 token# Worker 侧把 token 解包到 worker Subject topology.auto-credentials: - org.apache.storm.hdfs.security.AutoHDFS - org.apache.storm.hbase.security.AutoHBase # Nimbus 侧代表提交者获取 token nimbus.autocredential.plugins.classes: - org.apache.storm.hdfs.security.AutoHDFSNimbus - org.apache.storm.hbase.security.AutoHBaseNimbus # Nimbus 侧为长跑拓扑续期 token nimbus.credential.renewers.classes: - org.apache.storm.hdfs.security.AutoHDFSNimbus - org.apache.storm.hbase.security.AutoHBaseNimbus相关的凭据参数包括hdfs.keytab.file/hdfs.kerberos.principalNimbus 获取 HDFS token 用的主体、hbase.keytab.file/hbase.kerberos.principal、topology.hdfs.uriNameNode URI默认取集群fs.defaultFS以及多集群场景下可选的hdfsCredentialsConfigKeys/hbaseCredentialsConfigKeys配置键列表。若只需其中一种仅保留对应的 HDFS 或 HBase 条目即可。完整的安全集群搭建Kerberos、impersonation、ACL见 docs/SECURITY.md。迁移与运维类工具让集群演进更平滑storm-blobstore-migration本地 Blobstore 到 HDFS Blobstore 迁移external/storm-blobstore-migration/README.md 提供了把 Nimbus 的 blob 从LocalFsBlobStore迁移到HdfsBlobStore的工具。工作流是make构建 tarball → 拷到 Nimbus 主机解压 → 编写名为config的配置文件必须含HDFS_BLOBSTORE_DIR、LOCAL_BLOBSTORE_DIR、HADOOP_CLASSPATH可选BLOBSTORE_PRINCIPAL、KEYTAB_FILE、JAAS_CONF→ 运行listHDFS.sh/listLocal.sh盘点两侧已有 blob或migrate.sh执行迁移。Nimbus 侧迁移的要点是先关闭所有 Nimbus 实例并备份配置再修改blobstore.dir、blobstore.hdfs.principal、blobstore.hdfs.keytab、blobstore.replication.factor、nimbus.blobstore.class五项配置确保STORM_EXT_CLASSPATH包含HADOOP_CLASSPATH然后在主 Nimbus 上运行migrate.sh核对listHDFS.sh/listLocal.sh结果后重启 Nimbus失败时恢复备份配置即可回退。Supervisor 侧则需关停 supervisor、配置supervisor.blobstore.class等项、杀掉残留 worker 进程并清空本地状态再重启README 说明这是因迁移过程中出现过只有清空本地状态才能解决的偶发错误。Hadoop jar 不随 Storm 或该工具打包需单独安装。storm-kafka-monitorKafka Spout 消费滞后查询external/storm-kafka-monitor/README.md 描述了一个查询 Kafka Spout 消费滞后lag并在 Storm UI 中展示的工具。它对比 spout 已成功消费的 offset 与 Kafka 中的最新 offset同时支持新旧两代 Kafka SpoutUI 侧在缺少该模块时会优雅降级不显示 lag 并仅记录一次提示日志。jar 打包在完整发行版的lib-tools/storm-kafka-monitor下lite 发行版可用$STORM_HOME/bin/storm-kafka-monitor-fetch安装后重启 UI。命令行用法$STORM_HOME_DIR/bin/storm-kafka-monitor脚本实际运行org.apache.storm.kafka.monitor.KafkaOffsetLagUtil支持参数参数必选说明-t/--topics是逗号分隔的 topic 列表-b/--bootstrap-brokers是broker 地址列表-g/--groupid是consumer group id-s/--security-protocol否安全协议-c/--consumer-config否consumer 配置文件路径storm-metrics 与 storm-metrics-prometheus指标上报external/storm-metrics-prometheus/README.md 说明该模块包含一个把集群指标推送到 Prometheus Pushgateway 的 reporterorg.apache.storm.metrics.prometheus.PrometheusPreparableReporter供 Prometheus 实例抓取。配置只需在storm.yaml中启用插件并设置推送端点storm.daemon.metrics.reporter.plugins: - org.apache.storm.metrics.prometheus.PrometheusPreparableReporter storm.daemon.metrics.reporter.interval.secs: 10 # Prometheus Pushgateway 配置 storm.daemon.metrics.reporter.plugin.prometheus.job: job_name storm.daemon.metrics.reporter.plugin.prometheus.endpoint: localhost:9091 storm.daemon.metrics.reporter.plugin.prometheus.scheme: http storm.daemon.metrics.reporter.plugin.prometheus.basic_auth_user: storm.daemon.metrics.reporter.plugin.prometheus.basic_auth_password: storm.daemon.metrics.reporter.plugin.prometheus.skip_tls_validation: false同时需把该 jar 及其 Prometheus 传递依赖放入 Storm 安装目录的/lib。集群指标的通用机制可参见 docs/ClusterMetrics.md 与 docs/metrics_v2.md。仓库中的 external/storm-metrics 是配套的指标基础模块。如何在自己的拓扑中使用 external 模块综合各模块 README使用 external 模块有三条主流路径可按部署形态选择Maven 依赖 打包进拓扑 jar推荐尤其适合连接器场景在拓扑工程中以${storm.version}引入对应 artifact并用maven-shade-plugin合并打包HDFS 连接器强制要求此方式否则会出现No FileSystem for scheme: hdfs。此时集群端无需任何额外安装。放入extlibworker 与守护进程均加载如 external/storm-iceberg/README.md 所示发行版提供$STORM_HOME/bin/storm-iceberg-fetch等 helper 脚本可从 Maven Central 解析模块及其运行时依赖到$STORM_HOME/extlib适合不想把大依赖树塞进拓扑 jar 的场景脚本支持--version/--dest参数并向 Maven 透传参数Maven 只需在运行脚本的机器上可用生成 jar 后可拷到各 worker 主机改动后需重新提交拓扑使 classpath 生效。放入extlib-daemon仅守护进程加载用于 Nimbus/Supervisor 侧功能如 storm-autocreds 这类运行在守护进程上的插件安装后需重启 Nimbus 与 Supervisors。需要留意发行版差异某些模块如 storm-autocreds、storm-kafka-monitor只随完整版发行包携带 jarlite 发行包仅附 README需用对应 helper 脚本自行安装。总结external 模块是 Apache Storm 生态中内核之外、生态之内的官方集成层它遵循与 Storm 同版本发布、由 Committer Sponsor 认领维护的治理模式覆盖 Kafka、JMS、Redis、JDBC、HDFS、Iceberg 等主流数据系统并提供凭据自动化、blobstore 迁移、offset 迁移与消费滞后监控等运维工具。从 external/pom.xml 的模块清单到每个模块的 README、配套示例examples与深度文档docs都可以在仓库中直接查阅与验证。选型时建议按数据接入 / 数据输出 / 安全凭据 / 迁移运维四类需求定位对应模块并始终遵循依赖版本与 Storm 一致、按部署形态选择打包路径、特殊连接器遵守 shade 打包要求三个原则。赞分享大数据流处理后端【免费下载链接】stormApache Storm项目地址https://gitcode.com/gh_mirrors/storm6/storm点击查看免费下载相关推荐StarRocks BE DataSink 模块架构解析六大数据汇模块与模块边界治理机制StarRocks BE DataSink 模块架构解析六大数据汇模块与模块边界治理机制 导读 本文聚焦 StarRocks 后端BE的数据汇DataS数据库OLAP数据仓库大数据湖仓一体数据分析ModelScope 多模态 Pipeline 全景解析multi_modal 模块 15 大 Pipeline 的架构与实战指南ModelScope 多模态 Pipeline 全景解析multi_modal 模块 15 大 Pipeline 的架构与实战指南 本文围绕 ModelSco人工智能大模型微调模型评测预训练OptiScaler 完整指南1 个 Insert 键切换 DLSS、FSR、XeSS免费给 DX12 游戏加帧生成OptiScaler 完整指南1 个 Insert 键切换 DLSS、FSR、XeSS免费给 DX12 游戏加帧生成 OptiScaler 是一款免费开源的图形学游戏开发上一篇如何在Windows电脑上安装APK文件APK安装器完整使用指南下一篇SRWE窗口编辑器打破Windows窗口限制的终极解决方案创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表