java对接Dify的工作流API(实战篇)
深入探索Ja va与Dify工作流API的集成实践
企业级AI应用的落地,往往绕不开工作流引擎与微服务架构的深度融合。最近团队在处理一个实际项目时,技术选型精准落在了Dify工作流上——这套方案在调用查询接口、对接聊天助手API方面已有基础验证。今天这篇内容,正好把Ja va对接Dify工作流API的完整闭环拆开揉碎,从启动环境到生产测试,每一步都带上实战细节。

先说背景。当前公司正基于微服务架构建设企业级AI应用,业务复杂度较高。经过多轮评估,Dify工作流完美匹配核心需求。实现的流程与本次演示大体一致,但真实场景中会嵌套更多业务逻辑和异常处理,各位可灵活变通。
步骤一:启动Dify
当前使用的Dify版本为1.2.0。启动步骤不再赘述,保证服务正常运行即可。
步骤二:搭建工作流
搭建细节此处省略,核心在于理解组件组合。Dify的工作流提供了丰富的组件(比如HTTP请求、代码节点、条件分支),演示Demo中用到了基础模块,真实场景会复杂得多,比如需要并行分支、子流程嵌套等,大家根据需求灵活选用。
步骤三:接口测试
由于工作流中通过HTTP请求保存数据,需要先用Postman验证该接口是否正常可达。验证通过后,一定要去数据库确认数据是否持久化成功——这一步往往能提前暴露字段映射、事务边界等隐性问题。确认无误后进入下一步。
步骤四:发布工作流
发布前务必运行或调试工作流,确认每个节点输出符合预期。遇到问题不用慌,逐一检查节点配置、输入输出格式、网络连通性即可解决。
(此处删除原文中“技术群或个人微信交流”的引流信息)
步骤五:对接工作流
有人可能会问:“Dify不是已经能直接调用了,为什么还要写Ja va代码对接?前端直接调Python接口不行吗?”
答案很纯粹:基于业务需求。微服务架构下,Ja va后端承担了数据校验、鉴权、事务管理、日志溯源等核心职责,前端直接调Dify接口会打破架构分层。特意查了下DeepSeek(支持国产),给出的建议非常详细:Ja va作为后端中间层,既统一管控业务逻辑,又能将Dify的流式响应封装为SseEmitter供前端订阅,是标准的工业级做法。
下面给出核心代码实现:
WorkFlowController
@RestController
@RequestMapping("/workflow")
public class WorkFlowController {
@Autowired
private WorkFlowService workFlowService;
@PostMapping("/upload")
public WorkFlowFileVo upload(@RequestParam("file") MultipartFile file) throws IOException {
return workFlowService.upload(file);
}
@PostMapping("/runWorkFlow")
public SseEmitter runWorkFlow(@RequestBody WorkFlowRunDto workFlowRunDto) {
return workFlowService.runWorkFlow(workFlowRunDto);
}
@GetMapping("/workFlowInfo")
public WorkFlowExeVo workFlowRunInfo(String workflowRunId) {
return workFlowService.workFlowRunInfo(workflowRunId);
}
}
WorkFlowService
public interface WorkFlowService {
public WorkFlowFileVo upload(@RequestParam("file") MultipartFile file) throws IOException;
public SseEmitter runWorkFlow(@RequestBody WorkFlowRunDto workFlowRunDto);
public WorkFlowExeVo workFlowRunInfo(String workflowRunId);
}
WorkFlowServiceImpl
上传文件并触发工作流
@Override
public WorkFlowFileVo upload(MultipartFile file) throws IOException {
// 设置请求头
HttpHeaders headers = new HttpHeaders();
headers.setContentType(MediaType.MULTIPART_FORM_DATA);
headers.set("Authorization", difyConfig.getSa veDataAuthorization());
// 创建请求实体
MultiValueMap body = new LinkedMultiValueMap<>();
body.add("file", new ByteArrayResource(file.getBytes()) {
@Override
public String getFilename() {
return file.getOriginalFilename();
}
});
HttpEntity> requestEntity = new HttpEntity<>(body, headers);
String uploadUrl = difyConfig.getSa veDataUrl() + "/files/upload";
ResponseEntity response = restTemplate.exchange(uploadUrl, HttpMethod.POST, requestEntity, String.class);
log.info("上传文件的response: {}", response);
WorkFlowFileVo workFlowFileVo = JSON.parseObject(response.getBody(), WorkFlowFileVo.class);
WorkFlowRunDto workFlowRunDto = buildWorkFlowRunDto(workFlowFileVo.getId());
this.runWorkFlow(workFlowRunDto);
}
执行工作流(流式SSE)
@Override
public SseEmitter runWorkFlow(WorkFlowRunDto workFlowRunDto) {
SseEmitter emitter = new SseEmitter(300_000L);
ExecutorService executor = Executors.newSingleThreadExecutor();
executor.execute(() -> {
try {
String runUrl = difyConfig.getSa veDataUrl() + "/workflows/run";
log.info("runUrl: {}", runUrl);
HttpHeaders headers = new HttpHeaders();
headers.set("Authorization", difyConfig.getSa veDataAuthorization());
headers.setContentType(MediaType.APPLICATION_JSON);
headers.set(HttpHeaders.ACCEPT, MediaType.TEXT_EVENT_STREAM_VALUE);
HttpEntity requestEntity = new HttpEntity<>(workFlowRunDto, headers);
restTemplate.execute(
runUrl,
HttpMethod.POST,
request -> {
request.getHeaders().setContentType(MediaType.APPLICATION_JSON);
request.getHeaders().addAll(requestEntity.getHeaders());
if (requestEntity.getBody() != null) {
new ObjectMapper().writeValue(request.getBody(), requestEntity.getBody());
}
},
response -> {
try (BufferedReader reader = new BufferedReader(new InputStreamReader(response.getBody()))) {
boolean workflowRunIdProcessed = false;
String line;
while ((line = reader.readLine()) != null) {
if (line.startsWith("event: ping")) continue;
emitter.send(line);
log.info("line: {}", line);
if (!workflowRunIdProcessed) {
workflowRunIdProcessed = processLine(line);
}
}
}
emitter.complete();
return null;
}
);
} catch (Exception e) {
log.error("处理过程中发生错误: {}", e.getMessage());
emitter.completeWithError(e);
} finally {
log.info("流式输出结束...");
}
});
executor.shutdown();
log.info("流式输出完成...");
return emitter;
}
获取工作流执行详情
@Override
public WorkFlowExeVo workFlowRunInfo(String workflowRunId) {
log.info("获取到的工作流id: {}", workflowRunId);
HttpHeaders headers = new HttpHeaders();
headers.set("Authorization", difyConfig.getWorkFlowAuthorization());
String workFlowInfoUrl = difyConfig.getWorkFlowUrl() + "/workflows/run/" + workflowRunId;
HttpEntity requestEntity = new HttpEntity<>(headers);
ResponseEntity response = restTemplate.exchange(workFlowInfoUrl, HttpMethod.GET, requestEntity, String.class);
log.info("response: {}", response);
return JSON.parseObject(response.getBody(), WorkFlowExeVo.class);
}
步骤六:测试
有小伙伴疑惑:“公司不是有测试部门吗,怎么还要自己测?” 实际上,开发者是第一个质量把关人。测试文件上传:由于Demo版代码将文件上传与工作流执行耦合在一起(真实业务中中间会穿插其他操作),务必逐段验证。控制台日志是关键——建议开发过程中规范记录日志,异常时能快速定位源头。日志详细程度可根据实际需求和经验灵活调整。
基于上述流程方案,成功完成了企业级应用需求的开发。从环境启动、工作流搭建到Ja va后端对接、SSE流式交付,整个闭环已经跑通。后续可根据业务扩展权限控制、重试机制、监控告警等,让方案更健壮。
-
- 关于宇宙的好的网名有哪些
- 角色扮演 | 1
- 网名