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 { 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); } } }