本篇技术指南以 README.md 为骨架,结合当前仓库源码,系统讲解 Hazelcast 作为"统一实时数据平台"的核心定位、典型应用场景、关键能力(流处理引擎 Jet、分布式数据结构、SQL 查询、连接器生态)以及从源码构建与测试的完整方法。读完本文,你将掌握 Hazelcast 解决实时数据处理问题的整体架构思路,理解 Jet 引擎的 Pipeline 编程模型,并能在本地完成源码编译、运行测试与基础配置。
什么是 Hazelcast
Hazelcast 是一个统一的实时数据平台(unified real-time data platform),它将流处理(stream processing)与快速数据存储(fast data store)融合在同一套分布式系统之中。其核心理念是:在数据进入数据库或数据湖"落地"之前,就对"流动中的数据"(data in motion)进行处理、富化和即时响应,从而让业务在毫秒级延迟内做出反应。
README 的表述可以概括为:企业使用 Hazelcast 处理流式数据,结合历史上下文(存储在 Hazelcast 内存数据网格中的上下文数据)进行增强,并通过标准规则或ML/AI 驱动的自动化在数据落库之前即时采取行动,从而创造新的收入流、降低风险、提升运营效率。
从当前仓库的模块结构(pom.xml 中的<modules>)可以看到这个"平台"的构成,它远不止一个内存缓存:
| 模块 | 作用 |
|---|---|
| hazelcast | 核心引擎,包含分布式数据网格与 Jet 流处理引擎 |
| hazelcast-sql | SQL 查询引擎,基于 Apache Calcite 定制(见hazelcast-sql/src/main/java/org/apache/calcite) |
| hazelcast-vector | 向量数据存储与检索(hazelcast-vector/src/main/java/com/hazelcast/vector) |
| hazelcast-tpc-engine | 高性能线程池(Thread-Per-Core)执行引擎 |
| hazelcast-spring、hazelcast-spring-boot-autoconfiguration | 与 Spring / Spring Boot 生态集成 |
| extensions | 连接器与扩展库(Kafka、Hadoop、S3、CDC 等) |
| distribution | 发行版组装、启动脚本(hz、hz-cli等)与 JVM 参数模板 |
当前仓库版本为6.0.0-SNAPSHOT(见 pom.xml 的<version>声明),属于社区开源版(Hazelcast Community)。
何时使用 Hazelcast:典型应用场景
README 明确指出,Hazelcast 是一个能承载多种实时应用工作负载的平台。适合使用它的场景包括:
- 有状态的流式/批量数据处理:对流动中的流数据或静态存量数据做有状态处理;
- 直接用 SQL 查询流式与批量数据源:流和批统一到 SQL 查询体系中;
- 通过连接器库接入数据,并用低延迟 SQL 对外提供查询服务:连接器负责"摄入",低延迟 SQL 负责"服务";
- 事件驱动的应用推送:事件发生时即时向应用推送更新(事件监听、事件日志);
- 低延迟的队列式或发布/订阅(pub-sub)消息通信;
- 通过缓存模式(read/write-through、write-behind)快速访问上下文数据与事务数据;
- 微服务的分布式协调;
- 跨地域或同地域数据中心间的数据复制(WAN)。
这些场景对应的具体实现散落在 hazelcast/src/main/java/com/hazelcast 下的各个子包中:map(分布式 KV 存储)、topic(pub-sub)、collection(队列/列表/集合)、multimap、replicatedmap、ringbuffer、cp(CP 子系统,用于分布式协调)、wan(跨数据中心复制)、sql(SQL 引擎入口)等。
核心特性一览
README 列出的关键特性包括:
- 有状态、容错的数据处理与查询:对数据流和静态数据使用 SQL 或数据流 API(Dataflow API / Jet Pipeline API)进行查询与处理;
- 全面的连接器库:Kafka、Hadoop、S3、RDBMS(JDBC)、JMS 等,见 extensions 模块;
- 分布式消息:基于 pub-sub 主题 与 队列;
- 分布式、分区化、可查询的键值存储:带事件监听器,可作为事件流的上下文数据源以低延迟读取,即 IMap;
- 与 Python 的机器学习模型集成:可将 Python 训练的 ML 模型部署到数据处理流水线中(见 extensions/python);
- 云原生、随处可运行的架构;
- 滚动升级(rolling upgrade)实现零停机运维;
- 流处理管道的at-least-once 与 exactly-once 处理保障;
- 跨数据中心与地理区域的WAN 数据复制;
- 键值点查与 pub-sub 的微秒级性能;
- 客户端库覆盖Java、Python、Node.js、.NET、C++、Go等多种语言。
说明:README 中提到的具体性能数字(如单节点聚合 1000 万事件/秒、集群处理十亿事件/秒等)来自官方发布的基准与博客,属于官方对外声明,本文不加以验证,也不作为事实性承诺引用。
深入核心:Jet 流处理引擎与有状态数据处理
README 中单独强调的"Stateful Data Processing"部分指出:Hazelcast 内置了一个名为Jet的数据处理引擎,可用于构建流式/实时与批量/静态两类数据管道,且管道是**弹性(elastic)**的——即随集群成员增减自动伸缩。
Jet 的编程模型位于 hazelcast/src/main/java/com/hazelcast/jet 包中,核心抽象如下:
- Pipeline:分布式计算任务的建模方式,类比"相互连接的水管系统"。基本元素是stage(阶段)——每个阶段接收上游数据、变换后导向下游。
Pipeline.create()创建一个空管道,readFrom(BatchSource)接入批量数据源,readFrom(StreamSource)接入流式数据源。一个 Pipeline 中的 stage 与底层 DAG 顶点并非一一对应:某些 stage 会展开为多个顶点(如分组操作是级联的两个顶点),某些 stage 级联(如连续的 map/filter/flatMap)则会被融合进单个顶点以提升性能。 - Sources / Sinks:声明式地描述数据源与数据汇,包括 Map、Cache、List、文件、Socket、JDBC、JMS 等内建来源;Sinks 则对应写出到 Map、Cache、List、Socket、文件、JDBC、JMS、Observable 等目标。
- DAG / Processor:更底层的 Core API(
com.hazelcast.jet.core),提供DAG(有向无环图)、Vertex、Edge(支持 partitioned、broadcast、allToOne 等路由策略)、Processor、ProcessorSupplier等原语,Pipeline 最终会翻译成 DAG 执行。 - 窗口与聚合:
WindowDefinition(滑动/翻滚/会话窗口)、AggregateOperations、GroupAggregateBuilder等提供事件时间窗口聚合能力。 - 容错保障:
ProcessingGuarantee配合状态快照(snapshot)实现 at-least-once / exactly-once;Job、JobStatus、JobStatusListener用于任务的提交、状态跟踪与监听。
连接器扩展(如 Kafka)即通过实现这些底层 Processor 接口接入 Jet 的,例如 extensions/kafka 模块(com.hazelcast.jet.kafka包)提供了KafkaSources与KafkaSinks,让 Pipeline 可以像读写内建数据结构一样读写 Kafka 主题。
快速上手:安装与运行
README 指引读者参照官方 Getting Started Guide 安装启动。在当前仓库中,与"运行一个 Hazelcast 成员"直接相关的资产包括:
- 发行版启动脚本:发行版的统一入口脚本(
hz、hz-cli、hz-start、hz-stop、hz-healthcheck等); - JVM 参数模板:发行版自带的 JVM 参数文件;
- 日志配置:发行版自带的 log4j2 日志配置。
默认配置解读:发行版默认配置位于 hazelcast-assembly.yaml(YAML 版本)与 hazelcast-assembly.xml(XML 版本,二者等价)。几个关键默认值:
- 集群名:
cluster-name: dev(YAML)/<cluster-name>dev</cluster-name>(XML)。同一集群的所有成员必须配置相同集群名,客户端连接时也必须使用它; - 监听端口:
5701,port-count: 100(成员会在 5701~5801 范围内尝试绑定),auto-increment: true(端口被占用时自动递增尝试); - 发现机制(join):
auto-detection(自动探测)、multicast(组播,默认禁用)、tcp-ip(默认指向127.0.0.1)、aws/gcp/azure/kubernetes/eureka(云厂商发现,默认均禁用)。注释中特别说明:Hazelcast 只与使用同一发现机制的节点组建成集群; - 接口绑定:ZIP/TAR 发行版默认只绑定回环地址
127.0.0.1(通过hazelcast.socket.bind.any属性控制),Docker 镜像则监听所有接口; - REST 端点:
rest-api默认开启,其中HEALTH_CHECK组开启、CLUSTER_READ组关闭; - Map 存储格式:默认
in-memory-format: BINARY(键值以二进制存储),可选 OBJECT、NATIVE; - 执行器线程池:
executor-service/scheduled-executor-service默认pool-size: 16;durable-executor 默认capacity: 100、durability: 1(每个任务一份主副本加一份备份副本)。
以上配置文件均可直接复制后按需修改,作为自定义集群配置的起点。
从源码构建 Hazelcast
README 给出了完整的源码构建流程。构建环境最低要求 JDK 17。
# 拉取最新代码 $ git pull origin master # 使用 Maven wrapper 构建(推荐),跳过测试 $ ./mvnw clean package -DskipTests仓库提供了 Maven wrapper 脚本(./mvnw),建议直接使用;也可以使用与本仓库 wrapper 脚本相同版本的本地 Maven 发行版执行同样的命令。
快速构建模式:-Dquick
除了完整构建,仓库还提供了quick 构建模式:设置-Dquick系统属性后,构建会跳过校验类任务(测试、Checkstyle 校验、Javadoc、source 插件等),并且不构建extensions与distribution模块,适合日常快速迭代本地开发。
这一点在根 pom.xml 中有明确对应实现:not-quickprofile 默认激活并纳入extensions、distribution、hazelcast-it三个模块,而设置-Dquick后会禁用该 profile,从而跳过这些附加模块。
构建产物与校验
- 根 pom.xml 将核心模块、Spring 集成、SQL 引擎、向量模块等组织为多模块 Maven 工程,顶层
hazelcast-root使用hazelcast-parent作为父 POM(hazelcast-parent/pom.xml); - 代码风格与许可证头校验由 checkstyle/checkstyle.xml 及 checkstyle/suppressions.xml 定义,新代码需符合 Apache 2.0 头(见 checkstyle/ClassHeaderApache.txt)。
测试体系:三种测试 Profile 与并行测试
README 特别提醒:默认构建会执行数千个测试,可能耗时相当长。Hazelcast 将测试分为三个 profile:
| Profile | 命令 | 用途 |
|---|---|---|
| 默认 | ./mvnw test | 快速/集成测试,可用-P parallelTest并行执行(无需网络) |
| 慢速测试 | ./mvnw test -P nightly-build | 执行较慢或无法并行运行的测试 |
| 全部测试 | ./mvnw test -P all-tests | 串行执行全部测试(需要网络) |
这些 profile 在根 pom.xml 中有完整定义,其行为可从配置细节中得到印证:
- 默认构建(surefire)默认排除
SlowTest与NightlyTest两组注解标记的测试,且以-Dhazelcast.test.use.network=false关闭网络; parallelTestprofile:以 CPU 核数一半(0.5C)的 fork 数量并行执行标注为ParallelJVMTest的测试,同时另设一个singlejvmexecution 执行其余测试;-Dhazelcast.test.multiple.jvm=true表明其支持多 JVM 测试模式;nightly-buildprofile:只执行NightlyTest与SlowTest两组测试(surefiregroups+ failsafe),用于定期长跑测试;all-testsprofile:禁用并行(parallel: none),串行执行所有测试,failsafe 以-Dhazelcast.test.use.network=true运行集成测试(*IT.java)。
此外还有两点对本地开发的实用建议:
- 部分测试依赖 Docker运行;可通过设置
-Dhazelcast.disable.docker.tests系统属性来忽略这些测试; - 开发 PR 时,本地只需运行新增测试及少量相关子集即可,完整测试套件由 CI 的 PR builder 负责运行。
许可证与贡献
本仓库源码遵循两种许可证之一(详见 LICENSE 与 licenses 目录):
- Apache License 2.0:仓库默认许可证;
- Hazelcast Community License:仅当文件头部明确标注时适用。
源码文件的许可证头模板参见 checkstyle/ClassHeaderApache.txt。贡献流程与开发约定见 CONTRIBUTING.md,安全相关的上报渠道见 SECURITY.md。
总结
Hazelcast 的核心价值在于把流处理引擎(Jet)、分布式内存数据网格与SQL 查询能力统一到一个平台上:事件流可以在内存中被即时处理、富化并驱动业务动作,同时上下文数据以低延迟就近读取。本文基于 README.md 梳理了其平台定位、适用场景与关键特性,并结合仓库源码深入讲解了 Jet 的 Pipeline 编程模型、发行版默认配置(集群名、端口、发现机制等)、JDK 17 + Maven wrapper 的源码构建方式,以及覆盖快速/慢速/全量三档的测试体系。开发者可以此为基础,进一步阅读 hazelcast/src/main/java/com/hazelcast/jet 下的引擎源码、extensions 中的连接器实现,以及 hazelcast-sql 中的 SQL 引擎,逐步掌握在 Hazelcast 上构建实时数据处理应用的方法。
- 缓存
- KV存储
- 消息队列
- 流处理
- 后端
【免费下载链接】hazelcast
Hazelcast is a unified real-time data platform combining stream processing with a fast data store, allowing customers to act instantly on>项目地址:https://gitcode.com/gh_mirrors/ha/hazelcast
相关推荐
如何使用VapeLabs Auto Bot:从安装到多账户管理的完整教程
如何使用VapeLabs Auto Bot:从安装到多账户管理的完整教程 VapeLabs Auto Bot是一款针对TheVapeLabs空投平台的自动化工具
Hazelcast与Kafka集成实战:构建企业级实时数据处理平台
Hazelcast与Kafka集成实战:构建企业级实时数据处理平台 在数字化转型浪潮中,企业对实时数据处理能力的需求日益迫切。传统批处理模式已无法满足业务对即时
缓存KV存储消息队列流处理后端实时数据处理新范式:Apache Airflow与流数据平台集成指南
实时数据处理新范式:Apache Airflow与流数据平台集成指南 引言:实时数据处理的痛点与解决方案 你是否还在为批处理任务无法满足实时数据需求而烦恼?是否
后端任务调度工作流自动化数据编排批处理数据工程流程编排