malei05 hace 1 mes
padre
commit
d3bcca62c2

+ 11 - 29
fn-rabbitmq/src/main/java/net/yaoyi/pipeline/rabbitmq/RabbitConsumer.java

@@ -22,8 +22,6 @@ public class RabbitConsumer implements PojoRequestHandler<Map<String, Object>, V
 
     private static final String PIPELINE_API = "http://1861991840278303.eventbridge.cn-beijing.aliyuncs.com/webhook/putEvents?token=62ed79329bb74d0d8ede5cc5d0487a45897a6ff99f5d4d61946daf9f6e9c7fbc766b9d55ef2241edb2c2cb3aee651ac101409d2783f341aa8e8d27fdb9869d86";
 
-    private static Connection connection;
-    private static Channel channel;
     private static final OkHttpClient httpClient = new OkHttpClient.Builder()
             .connectTimeout(10, TimeUnit.SECONDS)
             .readTimeout(30, TimeUnit.SECONDS)
@@ -35,45 +33,29 @@ public class RabbitConsumer implements PojoRequestHandler<Map<String, Object>, V
     public Void handleRequest(Map<String, Object> input, Context context) {
         log.info("函数计算启动, Request ID: {}", context.getRequestId());
 
-        try {
-            // 初始化 RabbitMQ 连接
-            initRabbitMQ();
-            // 启动消费者
-            startConsumer();
+        try (Connection connection = RabbitMQHelper.createConnection();
+             Channel channel = RabbitMQHelper.createChannel(connection, QUEUE_NAME)) {
+
+            startConsumer(channel);
             log.info("RabbitMQ 消费者启动成功");
-        } catch (Exception e) {
-            log.error("启动 RabbitMQ 消费者失败", e);
-            throw new RuntimeException("启动失败: " + e.getMessage(), e);
-        }
 
-        try {
-            if (input.containsKey("waitSeconds")) {
-                Thread.sleep((int) input.get("waitSeconds") * 1000);
+            Object val = input.get("waitSeconds");
+            if (val instanceof Integer) {
+                Thread.sleep(((int) val) * 1000);
             } else {
                 Thread.sleep(2 * 60 * 1000);
             }
-        } catch (InterruptedException e) {
-            Thread.currentThread().interrupt();
+        } catch (Exception e) {
+            log.error("启动 RabbitMQ 消费者失败", e);
+            throw new RuntimeException("启动失败: " + e.getMessage(), e);
         }
-
-        RabbitMQHelper.close(channel, connection);
         return null;
     }
 
-    /**
-     * 初始化 RabbitMQ 连接
-     */
-    private void initRabbitMQ() {
-        if (connection == null) {
-            connection = RabbitMQHelper.createConnection();
-            channel = RabbitMQHelper.createChannel(connection, QUEUE_NAME);
-        }
-    }
-
     /**
      * 启动消费者
      */
-    private void startConsumer() throws IOException {
+    private void startConsumer(Channel channel) throws IOException {
         channel.basicQos(1);
 
         log.info("开始消费队列: {}", QUEUE_NAME);