malei05 1 hónapja
szülő
commit
626a837780

+ 3 - 3
fn-common/src/main/java/net/yaoyi/pipeline/common/JacksonUtils.java

@@ -58,7 +58,7 @@ public class JacksonUtils {
         }
     }
 
-    public static String toJsonString(Object object) {
+    public static String toJson(Object object) {
         try {
             return objectMapper.writeValueAsString(object);
         } catch (Exception e) {
@@ -66,7 +66,7 @@ public class JacksonUtils {
         }
     }
 
-    public static <T> T fromJSON(String content, Class<T> valueType) {
+    public static <T> T fromJson(String content, Class<T> valueType) {
         try {
             return objectMapper.readValue(content, valueType);
         } catch (Exception e) {
@@ -74,7 +74,7 @@ public class JacksonUtils {
         }
     }
 
-    public static <T> List<T> fromJSONArray(String contend, Class<T> type) throws JsonProcessingException {
+    public static <T> List<T> fromJsonArray(String contend, Class<T> type) throws JsonProcessingException {
         return objectMapper.readValue(contend, objectMapper.getTypeFactory().constructCollectionType(List.class, type));
     }
 

+ 3 - 3
fn-opensearch/src/main/java/net/yaoyi/pipeline/opensearch/OpenSearchHandler.java

@@ -51,7 +51,7 @@ public class OpenSearchHandler implements PojoRequestHandler<TaskImageMsg, List<
             String nsTaskFilter = String.format("sourceSendPackageDeptId = %d and taskId = %d and id != '%s'", taskImageIndex.getSourceSendPackageDeptId(), taskImageIndex.getTaskId(), taskImageIndex.getId());
             List<TaskImageIndex.MultiResult> nsTaskResult = knnSearch(namespace, taskImageIndex.getVector(), nsTaskFilter, "task");
             multiResultList.addAll(nsTaskResult);
-            taskImageIndex.setMultiResult(JacksonUtils.toJsonString(multiResultList));
+            taskImageIndex.setMultiResult(JacksonUtils.toJson(multiResultList));
         }
 
         //先落库
@@ -114,11 +114,11 @@ public class OpenSearchHandler implements PojoRequestHandler<TaskImageMsg, List<
             JsonNode bodyNode = JacksonUtils.toJsonNode(responseBody);
 
             if (!bodyNode.has("status") || !"OK".equals(bodyNode.get("status").asText())) {
-                log.error("push document with multiResult failed, request:{}, response:{}", JacksonUtils.toJsonString(request), responseBody);
+                log.error("push document with multiResult failed, request:{}, response:{}", JacksonUtils.toJson(request), responseBody);
                 throw new RuntimeException("push document failed");
             }
         } catch (Exception e) {
-            log.error("push document with multiResult exception, request:{}", JacksonUtils.toJsonString(request), e);
+            log.error("push document with multiResult exception, request:{}", JacksonUtils.toJson(request), e);
             throw new RuntimeException(e);
         }
     }

+ 6 - 0
fn-rabbitmq/pom.xml

@@ -49,6 +49,12 @@
             <artifactId>okhttp</artifactId>
             <version>4.12.0</version>
         </dependency>
+        <dependency>
+            <groupId>net.yaoyi.pipeline</groupId>
+            <artifactId>fn-common</artifactId>
+            <version>1.0.0-SNAPSHOT</version>
+            <scope>compile</scope>
+        </dependency>
     </dependencies>
 
     <build>

+ 63 - 0
fn-rabbitmq/src/main/java/net/yaoyi/pipeline/rabbitmq/MultiResultProducer.java

@@ -0,0 +1,63 @@
+package net.yaoyi.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 lombok.extern.slf4j.Slf4j;
+import net.yaoyi.pipeline.common.JacksonUtils;
+import net.yaoyi.pipeline.model.AlgoBaseInfo;
+import net.yaoyi.pipeline.model.FnResult;
+
+import java.io.IOException;
+import java.nio.charset.StandardCharsets;
+import java.util.List;
+
+@Slf4j
+public class MultiResultProducer implements PojoRequestHandler<List<AlgoBaseInfo>, FnResult<?>> {
+
+    private static final String QUEUE_NAME = "vector.search.task.img.dul.check.result.queue";
+    private static final String EXCHANGE_NAME = "default.direct.exchange";
+
+    private Connection connection;
+    private Channel channel;
+
+    @Override
+    public FnResult<?> handleRequest(List<AlgoBaseInfo> algoBaseInfos, Context context) {
+        if (algoBaseInfos == null || algoBaseInfos.isEmpty()) {
+            return FnResult.ok();
+        }
+
+        //1.初始化RabbitMQ生产端
+        initRabbitMQ();
+
+        //2.执行消息发送
+        for (AlgoBaseInfo algoBaseInfo : algoBaseInfos) {
+            rabbitSend(JacksonUtils.toJson(algoBaseInfo));
+        }
+
+        //3.关闭生产端
+        closeRabbitMQ();
+        return FnResult.ok();
+    }
+
+    private void initRabbitMQ() {
+        connection = RabbitMQHelper.createConnection();
+        channel = RabbitMQHelper.createChannel(connection, QUEUE_NAME);
+    }
+
+    private void rabbitSend(String json) {
+        try {
+            channel.basicPublish(EXCHANGE_NAME, QUEUE_NAME, null, json.getBytes(StandardCharsets.UTF_8));
+            log.info("消息发送成功: {}", json);
+        } catch (IOException e) {
+            log.error("消息发送失败: {}", json, e);
+            throw new RuntimeException("消息发送失败: " + e.getMessage(), e);
+        }
+    }
+
+    private void closeRabbitMQ() {
+        RabbitMQHelper.close(channel, connection);
+    }
+
+}

+ 23 - 45
fn-rabbitmq/src/main/java/net/yaoyi/pipeline/rabbitmq/RabbitConsumer.java

@@ -4,8 +4,6 @@ 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.*;
 
@@ -13,24 +11,17 @@ import java.io.IOException;
 import java.nio.charset.StandardCharsets;
 import java.util.Map;
 import java.util.concurrent.TimeUnit;
-import java.util.concurrent.TimeoutException;
 
 /**
  * RabbitMQ 消费者 - 阿里云函数计算实现
  */
 @Slf4j
-public class RabbitConsumer implements PojoRequestHandler<Map<String,Object>, Void> {
+public class RabbitConsumer implements PojoRequestHandler<Map<String, Object>, 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()
@@ -41,7 +32,7 @@ public class RabbitConsumer implements PojoRequestHandler<Map<String,Object>, Vo
     private static final MediaType JSON = MediaType.get("application/json; charset=utf-8");
 
     @Override
-    public Void handleRequest(Map<String,Object> input, Context context) {
+    public Void handleRequest(Map<String, Object> input, Context context) {
         log.info("函数计算启动, Request ID: {}", context.getRequestId());
 
         try {
@@ -60,28 +51,18 @@ public class RabbitConsumer implements PojoRequestHandler<Map<String,Object>, Vo
         } catch (InterruptedException e) {
             Thread.currentThread().interrupt();
         }
+
+        RabbitMQHelper.close(channel, connection);
         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 initRabbitMQ() {
+        if (connection == null) {
+            connection = RabbitMQHelper.createConnection();
+            channel = RabbitMQHelper.createChannel(connection, QUEUE_NAME);
         }
     }
 
@@ -92,24 +73,21 @@ public class RabbitConsumer implements PojoRequestHandler<Map<String,Object>, Vo
         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);
-        });
+        channel.basicConsume(QUEUE_NAME, false, (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("消息已重新入队");
+                    }
+                },
+                consumerTag -> log.info("消费者被取消: {}", consumerTag));
     }
 
     private void processMessage(String message) {

+ 19 - 0
yyc-pipeline-model/src/main/java/net/yaoyi/pipeline/model/FnResult.java

@@ -0,0 +1,19 @@
+package net.yaoyi.pipeline.model;
+
+import lombok.Data;
+
+@Data
+public class FnResult<T> {
+    private int code;
+    private String msg;
+    private T data;
+
+    public boolean isOk() {
+        return code == 0;
+    }
+
+    public static FnResult<?> ok() {
+        return new FnResult<>();
+    }
+
+}