Spring Data、Spark Streaming 与 Flink 集成 Elasticsearch 实战
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");