You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何用纯Java实现Kafka消息转存至Docker部署的Open Search?

纯Java环境下Kafka消费者将消息存入OpenSearch的实现方案

前提说明

已完成Kafka生产者/消费者功能验证,Docker中OpenSearch已启动并创建目标索引(假设索引名为kafka_messages)。

步骤1:添加OpenSearch Java客户端依赖

在项目构建文件(如Maven的pom.xml)中添加与OpenSearch版本匹配的官方客户端依赖:

<dependency>
    <groupId>org.opensearch.client</groupId>
    <artifactId>opensearch-rest-client</artifactId>
    <version>2.11.0</version> <!-- 版本需与你的OpenSearch实例一致 -->
</dependency>
<dependency>
    <groupId>org.opensearch.client</groupId>
    <artifactId>opensearch-java</artifactId>
    <version>2.11.0</version>
</dependency>

步骤2:封装OpenSearch客户端工具类

创建工具类统一管理OpenSearch连接,避免重复初始化:

import org.opensearch.client.RestClient;
import org.opensearch.client.json.JsonpMapper;
import org.opensearch.client.json.jackson.JacksonJsonpMapper;
import org.opensearch.client.opensearch.OpenSearchClient;
import org.opensearch.client.transport.rest_client.RestClientTransport;
import java.io.IOException;

public class OpenSearchClientUtil {
    private static OpenSearchClient client;

    static {
        // 连接Docker中运行的OpenSearch(默认端口9200)
        RestClient restClient = RestClient.builder(
                new org.apache.http.HttpHost("localhost", 9200, "http")
                // 若开启认证,添加以下配置
                // .setHttpClientConfigCallback(httpClientBuilder -> 
                //     httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider)
                // )
        ).build();

        JsonpMapper jsonpMapper = new JacksonJsonpMapper();
        RestClientTransport transport = new RestClientTransport(restClient, jsonpMapper);
        client = new OpenSearchClient(transport);
    }

    public static OpenSearchClient getClient() {
        return client;
    }

    // 程序结束时关闭客户端
    public static void closeClient() throws IOException {
        if (client != null) {
            client._transport().close();
        }
    }
}

步骤3:修改Kafka消费者,添加消息入库逻辑

在消费者的消息处理循环中,调用OpenSearch客户端将消息写入指定索引:

import org.opensearch.client.opensearch.core.IndexRequest;
import org.opensearch.client.opensearch.core.IndexResponse;
import java.io.IOException;
import java.util.Arrays;
import java.util.Properties;
import org.apache.kafka.clients.consumer.ConsumerConfig;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.consumer.ConsumerRecords;
import org.apache.kafka.clients.consumer.KafkaConsumer;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.time.Duration;

public class KafkaConsumerMain {
    public static void main(String[] args) {
        String servers = "localhost:9092";
        String groupId = "id";
        String topic = "mytopic";
        String osIndex = "kafka_messages"; // 替换为你的OpenSearch索引名

        Properties properties = new Properties();
        properties.setProperty(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, servers);
        properties.setProperty(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        properties.setProperty(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        properties.setProperty(ConsumerConfig.GROUP_ID_CONFIG, groupId);
        properties.setProperty(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");

        KafkaConsumer<String, String> kafkaConsumer = new KafkaConsumer<>(properties);
        kafkaConsumer.subscribe(Arrays.asList(topic));

        Logger logger = LoggerFactory.getLogger(KafkaConsumerMain.class.getName());
        OpenSearchClient osClient = OpenSearchClientUtil.getClient();

        try {
            while (true) {
                ConsumerRecords<String, String> records = kafkaConsumer.poll(Duration.ofMillis(3000));
                for (ConsumerRecord<String, String> record : records) {
                    logger.info("Key:" + record.key() + ", Value:" + record.value());
                    logger.info("Partition:" + record.partition() + ", Offset:" + record.offset());

                    // 将消息写入OpenSearch
                    try {
                        IndexRequest<String> request = new IndexRequest.Builder<String>()
                                .index(osIndex)
                                .document(record.value())
                                .build();

                        IndexResponse response = osClient.index(request);
                        logger.info("消息存入OpenSearch成功,文档ID: {}", response.id());
                    } catch (IOException e) {
                        logger.error("存入OpenSearch失败: {}", e.getMessage(), e);
                        // 生产环境可添加重试或死信队列逻辑
                    }
                }
            }
        } finally {
            kafkaConsumer.close();
            try {
                OpenSearchClientUtil.closeClient();
            } catch (IOException e) {
                logger.error("关闭OpenSearch客户端失败: {}", e.getMessage(), e);
            }
        }
    }
}

关键注意事项

  • 版本一致性:OpenSearch客户端版本必须与Docker中运行的OpenSearch版本完全匹配,否则会出现兼容性问题。
  • 认证配置:若OpenSearch开启了用户名密码认证,需添加CredentialsProvider配置:
    import org.apache.http.auth.AuthScope;
    import org.apache.http.auth.UsernamePasswordCredentials;
    import org.apache.http.client.CredentialsProvider;
    import org.apache.http.impl.client.BasicCredentialsProvider;
    
    // ...
    CredentialsProvider credentialsProvider = new BasicCredentialsProvider();
    credentialsProvider.setCredentials(AuthScope.ANY,
            new UsernamePasswordCredentials("admin", "admin")); // 替换为你的账号密码
    
    RestClient restClient = RestClient.builder(
            new org.apache.http.HttpHost("localhost", 9200, "http")
    ).setHttpClientConfigCallback(httpClientBuilder -> 
            httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider)
    ).build();
    
  • 消息结构化:如果消息是JSON格式,可直接作为文档存入;若为普通字符串,OpenSearch会自动将其存入_source字段,也可根据需求构造结构化文档后入库。

内容的提问来源于stack exchange,提问作者user20174493

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.18 20:05:28