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

Quarkus从Kafka拉取消息并向REST端点发送JSON负载的实现方法

可行实现方案梳理

针对你目前的Quarkus + MicroProfile Reactive Messaging + RxJava2技术栈需求,我整理了几个实用的方案,帮你把Kafka拉取的消息转发送到REST端点:


方案1:使用Quarkus反应式REST Client(最推荐)

Quarkus的MicroProfile REST Client原生支持反应式调用,和RxJava2无缝集成,代码简洁且符合规范。

步骤1:定义反应式REST Client接口

先创建一个REST客户端接口,指定目标端点的基础地址和请求方式:

import org.eclipse.microprofile.rest.client.inject.RegisterRestClient;
import io.reactivex.Flowable;
import javax.ws.rs.POST;
import javax.ws.rs.Path;
import javax.ws.rs.Consumes;
import javax.ws.rs.Produces;
import javax.ws.rs.core.MediaType;

@RegisterRestClient(baseUri = "https://your-target-api.com")
public interface TargetApiClient {
    @POST
    @Path("/submit-message")
    @Consumes(MediaType.APPLICATION_JSON)
    @Produces(MediaType.APPLICATION_JSON)
    Flowable<Void> sendMessage(CustomMessage message);
}

步骤2:修改原有消息处理方法

去掉@Outgoing和@Broadcast注解,直接在@Incoming方法里调用REST Client:

import org.eclipse.microprofile.rest.client.inject.RestClient;
import io.reactivex.Flowable;
import javax.enterprise.context.ApplicationScoped;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

@ApplicationScoped
public class MessageProcessor {
    private static final Logger LOGGER = LoggerFactory.getLogger(MessageProcessor.class);

    @RestClient
    TargetApiClient apiClient;

    @Incoming("pre-check")
    public Flowable<Void> publishToApi(CustomMessage customMessage) {
        LOGGER.info("Message received from topic = {}", customMessage);
        if (customMessage.ready) {
            return apiClient.sendMessage(customMessage)
                .doOnComplete(() -> LOGGER.info("Successfully sent message to REST endpoint"))
                .doOnError(error -> LOGGER.error("Failed to send message to REST endpoint", error));
        } else {
            LOGGER.info("Message not ready, skipping");
            return Flowable.empty();
        }
    }
}

方案优势

  • 完全贴合MicroProfile规范,Quarkus自动处理客户端的生命周期、连接池等
  • 支持通过application.properties配置超时、重试、认证等参数
  • RxJava2的Flowable可以直接对接,无需额外转换

方案2:使用Vert.x Web Client(更灵活的底层控制)

如果需要对HTTP请求做更精细的控制(比如自定义请求头、处理特殊响应码),可以用Quarkus内置的Vert.x Web Client,它和RxJava2的绑定非常完善。

修改后的处理方法示例

import io.vertx.reactivex.ext.web.client.WebClient;
import io.vertx.core.json.JsonObject;
import io.reactivex.Flowable;
import javax.enterprise.context.ApplicationScoped;
import javax.inject.Inject;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

@ApplicationScoped
public class MessageProcessor {
    private static final Logger LOGGER = LoggerFactory.getLogger(MessageProcessor.class);

    @Inject
    WebClient webClient;

    @Incoming("pre-check")
    public Flowable<Void> publishToApi(CustomMessage customMessage) {
        LOGGER.info("Message received from topic = {}", customMessage);
        if (customMessage.ready) {
            // 将CustomMessage转换为JSON(假设实体类已配置Jackson注解)
            JsonObject jsonPayload = JsonObject.mapFrom(customMessage);
            
            return webClient.post(443, "your-target-api.com", "/submit-message")
                .putHeader("Content-Type", "application/json")
                .putHeader("Authorization", "Bearer your-token") // 示例:添加认证头
                .rxSendJson(jsonPayload)
                .flatMapCompletable(response -> {
                    if (response.statusCode() >= 200 && response.statusCode() < 300) {
                        LOGGER.info("Message sent successfully, status code: {}", response.statusCode());
                        return io.reactivex.Completable.complete();
                    } else {
                        LOGGER.error("Failed to send message, status code: {}", response.statusCode());
                        return io.reactivex.Completable.error(new RuntimeException("API call failed with status " + response.statusCode()));
                    }
                })
                .toFlowable()
                .doOnError(error -> LOGGER.error("Error sending message to REST endpoint", error));
        } else {
            LOGGER.info("Message not ready, skipping");
            return Flowable.empty();
        }
    }
}

方案优势

  • 支持所有HTTP请求的精细配置,比如自定义超时、重试策略、请求头
  • 直接基于Quarkus底层的Vert.x,性能优异
  • RxJava2的Observable/Flowable可以和Vert.x的Rx API无缝转换

方案3:自定义SmallRye Outgoing连接器(扩展场景)

如果需要保留通道的灵活性(比如未来可能切换回Kafka或添加其他输出方式),可以自定义一个SmallRye的Outgoing连接器。不过这个方案相对复杂,适合有扩展需求的场景:

  1. 实现OutgoingConnectorFactory接口,在工厂类里处理HTTP发送逻辑
  2. 在application.properties中配置通道使用自定义连接器
  3. 原有代码可以保留@Outgoing注解,只需要修改通道配置

这个方案适合需要标准化输出逻辑、复用给多个服务的场景,但对于当前需求来说,前两个方案已经足够简洁高效。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.08 07:38:12