RabbitConsumer.java 4.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123
  1. package net.yao.pipeline.rabbitmq;
  2. import com.aliyun.fc.runtime.Context;
  3. import com.aliyun.fc.runtime.PojoRequestHandler;
  4. import com.rabbitmq.client.Channel;
  5. import com.rabbitmq.client.Connection;
  6. import com.rabbitmq.client.ConnectionFactory;
  7. import com.rabbitmq.client.DeliverCallback;
  8. import lombok.extern.slf4j.Slf4j;
  9. import okhttp3.*;
  10. import java.io.IOException;
  11. import java.nio.charset.StandardCharsets;
  12. import java.util.concurrent.TimeUnit;
  13. import java.util.concurrent.TimeoutException;
  14. /**
  15. * RabbitMQ 消费者 - 阿里云函数计算实现
  16. */
  17. @Slf4j
  18. public class RabbitConsumer implements PojoRequestHandler<String, Void> {
  19. private static final String RABBITMQ_HOST = "rabbitmq-cn-39y4twfud01-cn-beijing-vpc.mq.amqp.aliyuncs.com";
  20. private static final int RABBITMQ_PORT = 5672;
  21. private static final String RABBITMQ_USERNAME = "MjpyYWJiaXRtcS1jbi0zOXk0dHdmdWQwMTpMVEFJNXQ2SDR1OGFxb3A3cGViNkRSOFo=";
  22. private static final String RABBITMQ_PASSWORD = "MDdCQ0FFNjFEMTQyMEZDN0UyMThDNkM4MDgwMkI5QzgyNjAzQTdGQzoxNzgyMTk2NzQxNjcy";
  23. private static final String QUEUE_NAME = "vector.search.task.img.dul.check.queue";
  24. private static final String V_HOST = "test";
  25. private static final String PIPELINE_API = "http://1861991840278303.eventbridge.cn-beijing.aliyuncs.com/webhook/putEvents?token=63b99b24a3d545d880fa86f57c5ae56cae00c81d2fc74d058e9ace260515163132f6e27deeaa4bd4bc8e2a21cf6b6d6517b8d1e6d33c4de89f00c3524a54095b";
  26. private static ConnectionFactory factory;
  27. private static Connection connection;
  28. private static Channel channel;
  29. private static final OkHttpClient httpClient = new OkHttpClient.Builder()
  30. .connectTimeout(10, TimeUnit.SECONDS)
  31. .readTimeout(30, TimeUnit.SECONDS)
  32. .writeTimeout(30, TimeUnit.SECONDS)
  33. .build();
  34. private static final MediaType JSON = MediaType.get("application/json; charset=utf-8");
  35. @Override
  36. public Void handleRequest(String input, Context context) {
  37. log.info("函数计算启动, Request ID: {}", context.getRequestId());
  38. try {
  39. // 初始化 RabbitMQ 连接
  40. initRabbitMQ();
  41. // 启动消费者
  42. startConsumer();
  43. log.info("RabbitMQ 消费者启动成功");
  44. } catch (Exception e) {
  45. log.error("启动 RabbitMQ 消费者失败", e);
  46. throw new RuntimeException("启动失败: " + e.getMessage(), e);
  47. }
  48. return null;
  49. }
  50. /**
  51. * 初始化 RabbitMQ 连接
  52. */
  53. private void initRabbitMQ() throws IOException, TimeoutException {
  54. if (factory == null) {
  55. factory = new ConnectionFactory();
  56. factory.setHost(RABBITMQ_HOST);
  57. factory.setPort(RABBITMQ_PORT);
  58. factory.setUsername(RABBITMQ_USERNAME);
  59. factory.setPassword(RABBITMQ_PASSWORD);
  60. factory.setVirtualHost(V_HOST);
  61. log.info("连接 RabbitMQ: {}:{}", RABBITMQ_HOST, RABBITMQ_PORT);
  62. connection = factory.newConnection();
  63. channel = connection.createChannel();
  64. // 声明队列
  65. channel.queueDeclare(QUEUE_NAME, true, false, false, null);
  66. log.info("队列声明成功: {}", QUEUE_NAME);
  67. }
  68. }
  69. /**
  70. * 启动消费者
  71. */
  72. private void startConsumer() throws IOException {
  73. channel.basicQos(1);
  74. log.info("开始消费队列: {}", QUEUE_NAME);
  75. DeliverCallback deliverCallback = (consumerTag, delivery) -> {
  76. String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
  77. log.info("收到消息: {}", message);
  78. try {
  79. processMessage(message);
  80. channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
  81. log.info("消息处理成功并已确认");
  82. } catch (Exception e) {
  83. log.error("处理消息失败", e);
  84. channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, true);
  85. log.warn("消息已重新入队");
  86. }
  87. };
  88. channel.basicConsume(QUEUE_NAME, false, deliverCallback, consumerTag -> {
  89. log.info("消费者被取消: {}", consumerTag);
  90. });
  91. }
  92. private void processMessage(String message) {
  93. RequestBody body = RequestBody.create(message, JSON);
  94. Request request = new Request.Builder().url(PIPELINE_API).put(body).build();
  95. try (Response response = httpClient.newCall(request).execute()) {
  96. if (response.isSuccessful()) {
  97. String responseBody = response.body() != null ? response.body().string() : "";
  98. log.info("Pipline API 请求成功:{}", responseBody);
  99. } else {
  100. log.info("Pipline API 请求失败:{}", response.code());
  101. }
  102. } catch (Exception e) {
  103. log.error("Pipline API 请求异常", e);
  104. }
  105. }
  106. }