接入流式会话智能体

Maven依赖

		<!-- 阿里智能体SDK -->
		<dependency>
			<groupId>com.alibaba</groupId>
			<artifactId>dashscope-sdk-java</artifactId>
			<!-- 获取最新版本号 https://mvnrepository.com/artifact/com.alibaba/dashscope-sdk-java -->
			<version>2.21.8</version>
		</dependency>

Controller层

	@PostMapping(value = "/callDeepSeek", produces = MediaType.TEXT_EVENT_STREAM_VALUE)
	@ApiOperation(value = "接入DeepSeek会话", notes = "接入DeepSeek会话")
	public SseEmitter callDeepSeek(@RequestBody @Validated ZiWayXboDTO.CallDeepSeek param){
		return ziWayXboService.callDeepSeek(param);
	}

service层

	/**
	 * 接入DeepSeek会话
	 * @param param 参数
	 * @return 响应回答
	 */
	SseEmitter callDeepSeek(ZiWayXboDTO.CallDeepSeek param);

service实现层

	@Resource
	private ThreadPoolTaskExecutor threadPoolTaskExecutor;

	private static final Map<String, SseEmitter> EMITTERS = new ConcurrentHashMap<>();

	@Transactional(rollbackFor = Exception.class)
	@Override
	public SseEmitter callDeepSeek(ZiWayXboDTO.CallDeepSeek params) {
		if (Objects.isNull(params.getChatId())) {
			params.setChatId(this.createChat(params.getMessage()));
		}
		WorkdataUser user = SecurityUtils.getUser();

		// 3分钟超时
		SseEmitter emitter = new SseEmitter(180000L);
		EMITTERS.put(params.getChatId().toString(), emitter);

		StringBuilder responseText = new StringBuilder();
		StringBuilder thoughtText = new StringBuilder();

		// 使用原子布尔值标记完成状态,避免重复操作
		AtomicBoolean isCompleted = new AtomicBoolean(false);

		long snowflakeId = idGeneratorSnowFlake.snowflakeId();

		emitter.onCompletion(() -> {
			if (isCompleted.compareAndSet(false, true)) {
				EMITTERS.remove(params.getChatId().toString());
			}
		});

		emitter.onTimeout(() -> {
			if (isCompleted.compareAndSet(false, true)) {
				EMITTERS.remove(params.getChatId().toString());
			}
		});

		emitter.onError((error) -> {
			if (isCompleted.compareAndSet(false, true)) {
				EMITTERS.remove(params.getChatId().toString());
				log.error("SSE连接发生错误, chatId: {}", params.getChatId(), error);
			}
		});

		final String messageOne = "我想为员工提供健康午餐,有什么方案?";
		final String messageTwo = "有适合大型会议的团餐定制服务吗?";

		threadPoolTaskExecutor.execute(() -> {
			SseEmitter emitterThread = EMITTERS.get(params.getChatId().toString());

			// 检查emitter是否还存在
			if (emitterThread == null) {
				log.warn("SSE连接已关闭, chatId: {}", params.getChatId());
				return;
			}

			try {
				// 发送初始消息
				emitterThread.send(SseEmitter.event()
						.data("[sessionId]:" + params.getChatId()));

				emitterThread.send(SseEmitter.event()
						.data("[snowflakeId]:" + snowflakeId));

				emitterThread.send(SseEmitter.event()
						.data("[createTime]:" + DateUtil.formatDateTime(new Date())));

				// 处理预设消息
				if (Objects.equals(params.getMessage(), messageOne) || Objects.equals(params.getMessage(), "商务服务")) {
					String response = "我们提供专业企业订餐服务。请告知企业用餐人数、人均餐标(如25-30元/份)及口味偏好,我们将有客户经理一对一为您设计专属菜单;\n 请联系石经理:18320979581 ";
					// 商务合作联系方式
					String businessCooperation = this.getBusinessCooperation();

					responseText.append(response);
					if (StrUtil.isNotBlank(businessCooperation)) {
						responseText.append(businessCooperation);
					}

					emitterThread.send(SseEmitter.event().data("[reasoning]:该问题属于专业的商务服务内容,为保证服务质量,我需提供平台客服经理信息给用户,方便用户进一步联系客服经理对接。"));
					emitterThread.send(SseEmitter.event().data("[response]:" + responseText));
					emitterThread.complete();
					// 保存会话信息
					this.saveChatRecord(params, responseText, thoughtText, snowflakeId, user);
				}

				if (Objects.equals(params.getMessage(), messageTwo)) {
					String response = "专为会议、庆典提供团餐定制!请告知时间、人数、预算及主题(如茶歇、自助、盒饭),我们将有客户经理一对一为您设计专属菜单;\n 请联系石经理:18320979581 ";
					// 商务合作联系方式
					String businessCooperation = this.getBusinessCooperation();

					responseText.append(response);
					if (StrUtil.isNotBlank(businessCooperation)) {
						responseText.append(businessCooperation);
					}

					emitterThread.send(SseEmitter.event().data("[reasoning]:该问题属于专业的商务服务内容,为保证服务质量,我需提供平台客服经理信息给用户,方便用户进一步联系客服经理对接。"));
					emitterThread.send(SseEmitter.event().data("[response]:" + responseText));
					emitterThread.complete();
					// 保存会话信息
					this.saveChatRecord(params, responseText, thoughtText, snowflakeId, user);
				}

				String format = LocalDateTimeUtil.format(LocalDate.now(), "yyyy-MM-dd");
				if (Objects.equals(params.getMessage(), "今天有什么新品吗?")) {
					params.setMessage("今天" + format + "有什么新品吗?");
				}

				// 调用DashScope的SDK
				ApplicationParam param = ApplicationParam.builder()
						.apiKey("sk-5d674a53c6694c3fad05de36945561fa")
						.appId("5727e5eeb6f045f196ae5a036c43ceb0")
						.prompt(params.getMessage())
						.incrementalOutput(true)
						.hasThoughts(true)
						.build();

				if (Objects.equals(params.getMessage(), "今天" + format + "有什么新品吗?")) {
					params.setMessage("今天有什么新品吗?");
				}

				Application application = new Application();
				Flowable<ApplicationResult> result = application.streamCall(param);

				result.blockingForEach(data -> {
					// 检查连接是否还活跃
					if (isCompleted.get()) {
						return;
					}

					try {
						List<ApplicationOutput.Thought> thoughtList = data.getOutput().getThoughts();
						List<ApplicationOutput.Thought> reasoning = thoughtList.stream()
								.filter(v -> Objects.equals(v.getActionType(), "reasoning"))
								.collect(Collectors.toList());

						String thought = reasoning.get(0).getThought();
						String messageText = data.getOutput().getText();

						// 深度思考
						if (StrUtil.isNotBlank(thought)) {
							thoughtText.append("[reasoning]:" + thought);
							emitterThread.send(SseEmitter.event().data("[reasoning]:" + thought));
						}
						// 响应文本
						if (StrUtil.isNotBlank(messageText)) {
							responseText.append("[response]:" + messageText);
							emitterThread.send(SseEmitter.event().data("[response]:" + messageText));
						}
					} catch (IOException e) {
						log.warn("发送SSE消息失败,可能连接已关闭", e);
					}
				});

				// 安全完成
				if (!isCompleted.get()) {
					emitterThread.complete();
					// 保存会话信息
					this.saveChatRecord(params, responseText, thoughtText, snowflakeId, user);
				}

			} catch (IOException | ApiException | NoApiKeyException | InputRequiredException e) {
				if (!isCompleted.get()) {
					responseText.append("服务器异常, 请稍后再试");
					try {
						emitterThread.send(SseEmitter.event().data("[error]:服务器异常, 请稍后再试"));
					} catch (IOException ex) {
						log.warn("发送错误消息失败", ex);
					}
					emitterThread.completeWithError(e);
				}
				log.error("调用DeepSeek失败, message: [{}]", e.getMessage(), e);
			} finally {
				// 确保资源清理
				if (isCompleted.compareAndSet(false, true)) {
					EMITTERS.remove(params.getChatId().toString());
				}
			}
		});

		return emitter;
	}

	public Long createChat(String firstChat) {
		FormTemplateTable createChatTable = this.getTableByBaseType(ZiWayConstant.CREATE_CHAT, tenantId);
		Assert.notNull(createChatTable, ZiWayConstant.CREATE_CHAT + "表单模板不存在");
		WorkdataUser user = SecurityUtils.getUser();
		long id = idGeneratorSnowFlake.snowflakeId();
		Map<String, Object> createChatParamMap = new HashMap<>(7);
		createChatParamMap.put("id", id);
		createChatParamMap.put("user", user.getId());
		createChatParamMap.put("first_chat", firstChat);
		this.commonFieldInsert(createChatParamMap);
		formTableMapper.saveFormTable(createChatTable.getTableName(), createChatParamMap);
		return id;
	}
Logo

北京人形旗下天工造物具身智能开源社区,聚焦具身天工与慧思开物两大平台

更多推荐