| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123 |
- package net.yao.pipeline.rabbitmq;
- import com.aliyun.fc.runtime.Context;
- import com.aliyun.fc.runtime.PojoRequestHandler;
- import com.rabbitmq.client.Channel;
- import com.rabbitmq.client.Connection;
- import com.rabbitmq.client.ConnectionFactory;
- import com.rabbitmq.client.DeliverCallback;
- import lombok.extern.slf4j.Slf4j;
- import okhttp3.*;
- import java.io.IOException;
- import java.nio.charset.StandardCharsets;
- import java.util.concurrent.TimeUnit;
- import java.util.concurrent.TimeoutException;
- /**
- * RabbitMQ 消费者 - 阿里云函数计算实现
- */
- @Slf4j
- public class RabbitConsumer implements PojoRequestHandler<String, Void> {
- private static final String RABBITMQ_HOST = "rabbitmq-cn-39y4twfud01-cn-beijing-vpc.mq.amqp.aliyuncs.com";
- private static final int RABBITMQ_PORT = 5672;
- private static final String RABBITMQ_USERNAME = "MjpyYWJiaXRtcS1jbi0zOXk0dHdmdWQwMTpMVEFJNXQ2SDR1OGFxb3A3cGViNkRSOFo=";
- private static final String RABBITMQ_PASSWORD = "MDdCQ0FFNjFEMTQyMEZDN0UyMThDNkM4MDgwMkI5QzgyNjAzQTdGQzoxNzgyMTk2NzQxNjcy";
- private static final String QUEUE_NAME = "vector.search.task.img.dul.check.queue";
- private static final String V_HOST = "test";
- private static final String PIPELINE_API = "http://1861991840278303.eventbridge.cn-beijing.aliyuncs.com/webhook/putEvents?token=63b99b24a3d545d880fa86f57c5ae56cae00c81d2fc74d058e9ace260515163132f6e27deeaa4bd4bc8e2a21cf6b6d6517b8d1e6d33c4de89f00c3524a54095b";
- private static ConnectionFactory factory;
- private static Connection connection;
- private static Channel channel;
- private static final OkHttpClient httpClient = new OkHttpClient.Builder()
- .connectTimeout(10, TimeUnit.SECONDS)
- .readTimeout(30, TimeUnit.SECONDS)
- .writeTimeout(30, TimeUnit.SECONDS)
- .build();
- private static final MediaType JSON = MediaType.get("application/json; charset=utf-8");
- @Override
- public Void handleRequest(String input, Context context) {
- log.info("函数计算启动, Request ID: {}", context.getRequestId());
- try {
- // 初始化 RabbitMQ 连接
- initRabbitMQ();
- // 启动消费者
- startConsumer();
- log.info("RabbitMQ 消费者启动成功");
- } catch (Exception e) {
- log.error("启动 RabbitMQ 消费者失败", e);
- throw new RuntimeException("启动失败: " + e.getMessage(), e);
- }
- return null;
- }
- /**
- * 初始化 RabbitMQ 连接
- */
- private void initRabbitMQ() throws IOException, TimeoutException {
- if (factory == null) {
- factory = new ConnectionFactory();
- factory.setHost(RABBITMQ_HOST);
- factory.setPort(RABBITMQ_PORT);
- factory.setUsername(RABBITMQ_USERNAME);
- factory.setPassword(RABBITMQ_PASSWORD);
- factory.setVirtualHost(V_HOST);
- log.info("连接 RabbitMQ: {}:{}", RABBITMQ_HOST, RABBITMQ_PORT);
- connection = factory.newConnection();
- channel = connection.createChannel();
- // 声明队列
- channel.queueDeclare(QUEUE_NAME, true, false, false, null);
- log.info("队列声明成功: {}", QUEUE_NAME);
- }
- }
- /**
- * 启动消费者
- */
- private void startConsumer() throws IOException {
- channel.basicQos(1);
- log.info("开始消费队列: {}", QUEUE_NAME);
- DeliverCallback deliverCallback = (consumerTag, delivery) -> {
- String message = new String(delivery.getBody(), StandardCharsets.UTF_8);
- log.info("收到消息: {}", message);
- try {
- processMessage(message);
- channel.basicAck(delivery.getEnvelope().getDeliveryTag(), false);
- log.info("消息处理成功并已确认");
- } catch (Exception e) {
- log.error("处理消息失败", e);
- channel.basicNack(delivery.getEnvelope().getDeliveryTag(), false, true);
- log.warn("消息已重新入队");
- }
- };
- channel.basicConsume(QUEUE_NAME, false, deliverCallback, consumerTag -> {
- log.info("消费者被取消: {}", consumerTag);
- });
- }
- private void processMessage(String message) {
- RequestBody body = RequestBody.create(message, JSON);
- Request request = new Request.Builder().url(PIPELINE_API).put(body).build();
- try (Response response = httpClient.newCall(request).execute()) {
- if (response.isSuccessful()) {
- String responseBody = response.body() != null ? response.body().string() : "";
- log.info("Pipline API 请求成功:{}", responseBody);
- } else {
- log.info("Pipline API 请求失败:{}", response.code());
- }
- } catch (Exception e) {
- log.error("Pipline API 请求异常", e);
- }
- }
- }
|