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连接器。不过这个方案相对复杂,适合有扩展需求的场景:
- 实现
OutgoingConnectorFactory接口,在工厂类里处理HTTP发送逻辑 - 在
application.properties中配置通道使用自定义连接器 - 原有代码可以保留
@Outgoing注解,只需要修改通道配置
这个方案适合需要标准化输出逻辑、复用给多个服务的场景,但对于当前需求来说,前两个方案已经足够简洁高效。
内容的提问来源于stack exchange,提问作者animusdx
相关产品推荐
相关产品推荐

