当前位置:首页 > 技术 > 正文内容

Spring Data、Spark Streaming 与 Flink 集成 Elasticsearch 实战

访客 技术 2026年8月2日 2

Spring Data Elasticsearch 快速集成指南

Spring Data 是 Spring 家族中用于统一数据访问层的开源项目,支持关系型数据库、NoSQL 及搜索引擎。它通过抽象通用操作接口,极大简化了持久化代码的编写。其中,Spring Data Elasticsearch 模块封装了对 Elasticsearch 的客户端调用,开发者无需直接操作 REST API 即可完成索引管理、文档增删改查和复杂查询。

环境准备与依赖配置

使用 Spring Boot 构建项目时,建议选择 2.3.x 或更高版本以兼容 Elasticsearch 7.x 系列。在 pom.xml 中引入核心依赖:

<dependency>
    <groupId>org.springframework.boot</groupId>
    <artifactId>spring-boot-starter-data-elasticsearch</artifactId>
</dependency>

连接配置与客户端初始化

通过自定义配置类继承 AbstractElasticsearchConfiguration,构建基于 RestHighLevelClient 的连接实例。该方式替代已废弃的 TransportClient,符合未来版本演进方向。

@Configuration
@ConfigurationProperties(prefix = "es")
public class EsClientConfig extends AbstractElasticsearchConfiguration {

    private String host;
    private Integer port;

    @Override
    @Bean
    public RestHighLevelClient elasticsearchClient() {
        return new RestHighLevelClient(
            RestClient.builder(new HttpHost(host, port, "http"))
        );
    }
}

实体映射与 Repository 接口定义

使用注解描述文档结构,如 @Document 指定索引名,@Field 设置字段类型及分词策略。示例商品类如下:

@Data
@Document(indexName = "shopping")
public class Item {
    @Id
    private Long itemId;

    @Field(type = FieldType.Text, analyzer = "ik_max_word")
    private String name;

    @Field(type = FieldType.Keyword)
    private String category;

    @Field(type = FieldType.Double)
    private Double price;

    @Field(type = FieldType.Keyword, index = false)
    private String imgUrl;
}

DAO 层只需声明接口并继承 ElasticsearchRepository,即可获得完整的 CRUD 和分页排序能力:

public interface ItemRepository extends ElasticsearchRepository<Item, Long> {
}

基本操作测试

借助 JUnit 编写单元测试验证功能:

  • 保存文档itemRepository.save(item)
  • 分页检索:传入 PageRequest.of(page, size, Sort.by("price").descending())
  • 条件搜索:使用 QueryBuilders.termQuery("category", "手机") 构造查询
  • 删除索引:通过 ElasticsearchRestTemplate.deleteIndex(Item.class)

Spark Streaming 实时写入 Elasticsearch

Apache Spark Streaming 提供了微批处理模型,适用于高吞吐场景下的实时数据分析。结合 Elasticsearch 可实现日志聚合、监控告警等应用。

关键依赖添加

<dependency>
    <groupId>org.apache.spark</groupId>
    <artifactId>spark-streaming_2.12</artifactId>
    <version>3.0.0</version>
</dependency>
<dependency>
    <groupId>org.elasticsearch.client</groupId>
    <artifactId>elasticsearch-rest-high-level-client</artifactId>
    <version>7.8.0</version>
</dependency>

数据流处理逻辑

从 Socket 接收文本流,解析后构造 JSON 文档写入 ES:

JavaStreamingContext jssc = new JavaStreamingContext(conf, Durations.seconds(3));
JavaDStream<String> stream = jssc.socketTextStream("localhost", 9999);

stream.foreachRDD(rdd -> {
    rdd.foreach(record -> {
        String[] parts = record.split(" ", 2);
        Map<String, Object> doc = Collections.singletonMap("data", parts[1]);

        try (RestHighLevelClient client = new RestHighLevelClient(
                RestClient.builder(new HttpHost("localhost", 9200)))) {

            IndexRequest req = new IndexRequest("spark_stream")
                    .id(parts[0])
                    .source(doc, XContentType.JSON);
            client.index(req, RequestOptions.DEFAULT);
        }
    });
});

jssc.start();
jssc.awaitTermination();

Flink 流式数据同步到 Elasticsearch

Apache Flink 以其低延迟、精确一次语义(Exactly-Once)和强大的状态管理著称,适合对实时性和准确性要求更高的场景。

Maven 依赖设置

<dependency>
    <groupId>org.apache.flink</groupId>
    <artifactId>flink-connector-elasticsearch7_2.11</artifactId>
    <version>1.12.0</version>
</dependency>

使用 ElasticsearchSink 写入数据

构建 Sink 连接器并将流数据发送至指定索引:

final StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();

DataStreamSource<String> input = env.socketTextStream("localhost", 9999);

List<HttpHost> hosts = Collections.singletonList(new HttpHost("127.0.0.1", 9200, "http"));

ElasticsearchSink.Builder<String> esSinkBuilder = new ElasticsearchSink.Builder<>(
    hosts,
    (element, ctx, indexer) -> {
        Map<String, String> json = new HashMap<>();
        json.put("content", element);
        IndexRequest request = Requests.indexRequest()
            .index("flink_index")
            .source(json);
        indexer.add(request);
    }
);

// 设置每条记录立即刷新
esSinkBuilder.setBulkFlushMaxActions(1);

input.addSink(esSinkBuilder.build());

env.execute("Flink To Elasticsearch");

相关文章

Linux crontab 详解

1) crontab 是什么cron 是 Linux 的定时任务守护进程;crontab 是用来编辑/查看“按时间周期执行命令”的表(cron table)。常见两类:用户 crontab:每个用户一份(crontab -e 编辑)系统级 crontab / cron.d:可指定执行用户(/etc/crontab、/etc/cron.d/*)2) crontab 时间...

富文本里可以允许的 HTML 属性

一、所有标签默认允许的安全属性(极少)class        (可选)id           (通常建议禁用)title️ 注意:id 容易被滥用做锚点注入,很多系统直接禁用class 允许的话最好只允许固定前缀(如 editor-*)二、a 标签允许属性<a href="" t...

Mac 安装 Node.js 指南

方法一:通过官网安装包(最简单,适合初学者)如果你只是想快速安装并开始使用,这是最直接的方法。访问 Node.js 官网。页面会显示两个版本:LTS (Recommended For Most Users):长期支持版,最稳定。建议选这个。Current:最新特性版,包含最新功能但可能不够稳定。下载 .pkg 安装包并运行。按照安装向导点击“下一步”即可完成。方法二:使用 Homebrew 安装(...

Dom\HTML_NO_DEFAULT_NS 的副作用:自动加闭合标签

在使用Dom\HTMLDocument时,Dom\HTML_NO_DEFAULT_NS 将禁止在解析过程中设置元素的命名空间, 此设置是为了与DOMDocument向后兼容而存在的。当使用它时,已知的一个副作用就是:自动加闭合标签例如 </img> 为什么会这样?当你使用:Dom\HTML_NO_DEFAULT_NS文档会变成 无命名空间模式,此时内部更接近 XML...

Laravel 事件和监听器创建

在 Laravel 中,使用 Artisan 命令创建 Events(事件) 和 Listeners(监听器) 是非常高效的。你可以通过以下几种方式来实现:1. 手动创建单个 Event如果你只想创建一个事件类,可以使用 make:event 命令:Bashphp artisan make:event UserRegistered执行后,文件将生成在 app/Even...

自定义域名解析神器 dnsmasq

什么是 dnsmasq?dnsmasq 是一个轻量级、功能强大的网络服务工具,专为小型和中等规模网络设计。它是一个综合的网络基础设施解决方案[1]。dnsmasq 能做什么?功能说明应用场景DNS 转发与缓存将 DNS 查询转发到上游服务器(ISP、Google DNS 等),并在本地缓存结果加快 DNS 查询速度,减少外部 DNS 流量本地 DNS解析本地网络设备的主机名,无需编辑&n...

发表评论

访客

◎欢迎参与讨论,请在这里发表您的看法和观点。