HERA-V4 워크플로우 동작 구조와 개발 가이드¶
이 문서는 HERA-V4 워크플로우의 등록·실행 구조를 설명하고, 주니어 개발자가 간단한 워크플로우와 TaskHandler를 직접 등록하고 실행할 수 있도록 안내한다.
1. 개요¶
핵심 개념¶
HERA-V4 워크플로우는 DB에 저장된 파이프라인 정의를 읽어 Job을 만들고, 각 Task를 Kafka 이벤트로 순차 실행하는 비동기 구조다.
| 개념 | 역할 |
|---|---|
| Config/Pipeline | 무엇을 어떤 순서로 실행할지 정의 |
| Job | 워크플로우 1회 실행 인스턴스 |
| Task | Job 안의 개별 실행 단계 |
| Context | Task 사이에서 공유하는 입력·출력 데이터 |
| Ticket | Kafka로 전달되는 Job/Task 실행 메시지 |
| Coordinator | Job과 Task 실행 순서를 조정 |
| Worker | Task 상태 변경과 실행 처리 |
| TaskHandler | 실제 업무 로직 수행 |
| Kafka/Spring Integration | 등록·실행·완료·실패 이벤트 전달 |
전체 실행 흐름¶
flowchart TD
A[워크플로우 정의 등록] --> B[(wfw_config)]
B --> C[POST /api/v1/wfw/run]
C --> D[Coordinator.registerJob]
D --> E[(wfw_job 및 Context 저장)]
E --> F[JobRegistration 이벤트]
F --> G[Coordinator.startJob]
G --> H[JobExecution 이벤트]
H --> I[첫 번째 Task 등록]
I --> J[TaskRegistration 이벤트]
J --> K[Worker.startTask]
K --> L[TaskExecution 이벤트]
L --> M[TaskHandler 실행]
M -->|성공| N[TaskCompletion]
M -->|실패| O[TaskError]
N --> P{다음 Task 존재?}
P -->|예| I
P -->|아니요| Q[JobCompletion]
O --> R[Job 실패 처리]
Job, Task, ControlTask는 각각 registration, execution, completion, error 라이프사이클을 가지며 총 12개 Kafka 바인딩으로 연결된다.
2. 빠른 시작¶
워크플로우 등록¶
워크플로우 정의는 최종적으로 wfw_config에 저장된다.
주요 요청 필드는 다음과 같다.
{
"workflowId": "sample-workflow",
"workflowName": "샘플 워크플로우",
"trigger": "1",
"scheduler": "",
"activeYn": "Y",
"description": "워크플로우 예제",
"pipeline": "파이프라인 YAML 문자열"
}
workflowId를 생략하면 내부에서 ID가 생성된다. 등록 시 ConfigBuilder가 pipeline 문자열을 YAML로 파싱하고 Task가 하나 이상인지 검증한 후 저장한다. CRON, 고정 주기, 고정 일시 트리거는 스케줄 정보도 갱신한다.
화면에서는 hera-webapp의 워크플로우 설정 메뉴를 이용할 수 있다. 화면 요청은 Feign Client를 거쳐 bootstrap-workflow로 전달된다.
가장 단순한 파이프라인¶
inputs:
- name: yourName
type: string
label: 사용자 이름
required: true
pre:
- type: io/print
label: 시작 로그
text: "워크플로우를 시작합니다."
tasks:
- type: io/print
label: 인사 출력
text: "Hello ${yourName}"
- type: random/int
label: 임의 숫자 생성
output: randomNumber
startInclusive: 10
endInclusive: 100
- type: time/sleep
label: 잠시 대기
millis: 1000
timeout: 5M
- type: io/print
label: 결과 출력
text: "생성된 숫자는 ${randomNumber}입니다."
post:
- type: io/print
label: 후처리
text: "본 작업이 끝났습니다."
finalize:
- type: io/print
label: 최종 처리
text: "워크플로우를 종료합니다."
outputs:
- name: resultNumber
value: "${randomNumber}"
inputs: 실행 시 외부에서 받을 값pre: 본 Task 실행 전 수행할 보조 Tasktasks: 실제 업무 단계post: 본 Task가 정상 처리된 후 수행finalize: Task 실행 마무리 단계outputs: Job 종료 시 노출할 결과${변수명}: Context 값 참조output: Task 반환값을 Context에 저장할 이름
예를 들어 random/int가 42를 반환하고 output: randomNumber가 지정되어 있으면 Context에 randomNumber=42가 저장된다. 이후 Task에서 ${randomNumber}로 참조한다.
워크플로우 실행¶
curl -X POST 'http://localhost:18020/api/v1/wfw/run' \
--data-urlencode 'workflowId=sample-workflow' \
--data-urlencode 'yourName=hera'
현재 직접 실행 API는 JSON body가 아니라 request parameter 방식이다. API는 실행 완료를 기다리지 않고 200 OK를 반환하며 실제 처리는 Kafka를 통해 비동기로 이어진다.
실행 요청을 받으면 Coordinator.registerJob()이 JobTicket을 만들고 Job 이력과 Context를 저장한 뒤 JobRegistration 이벤트를 발행한다. 이어 Job 상태가 STARTED로 바뀌고 첫 번째 Task가 등록된다.
3. 단계별 적용 가이드¶
Spring Integration의 역할¶
HERA-V4 워크플로우를 단순히 “Kafka 기반 워크플로우”라고 이해하면 내부 실행 구조의 절반만 본 것이다. 현재 구현은 다음 세 계층을 조합한다.
| 계층 | 담당 기술 | 역할 |
|---|---|---|
| 외부 이벤트 운반 | Kafka | Job·Task 라이프사이클 메시지 저장 및 전달 |
| Kafka 바인딩 | Spring Cloud Stream | Kafka 토픽과 Java Consumer/Producer 연결 |
| 내부 오케스트레이션 | Spring Integration | Channel, Router, Handler, Filter, Gateway, Advice로 실행 경로 구성 |
flowchart LR
A[Coordinator / Worker] --> B[WorkflowMessageGateway]
B --> C[SI Out Channel]
C --> D[Outbound IntegrationFlow]
D --> E[StreamBridge]
E --> F[Spring Cloud Stream Binding]
F --> G[Kafka Topic]
G --> H[Spring Cloud Stream Consumer]
H --> I[SI DirectChannel]
I --> J[IntegrationFlow]
J --> K[Router / Handler / Filter]
K --> L[Coordinator / Worker / Executor]
Kafka는 이벤트를 다음 처리 사이클로 넘기는 비동기 경계다. Spring Integration은 Kafka에서 받은 메시지가 서비스 내부에서 어떤 업무 메서드로 이동할지를 결정한다.
주요 Spring Integration 구성요소¶
| 구성요소 | HERA-V4에서의 용도 |
|---|---|
Message<JobTicket> / Message<TaskTicket> | 실행 데이터와 메시지 Header 전달 |
DirectChannel | 서비스 내부의 동기 메시지 전달 |
IntegrationFlow | Channel과 Router·Handler·Filter 연결 |
route() | Job/Task 상태에 따른 실행 경로 선택 |
handle() | Coordinator, Worker, Executor 호출 |
filter() | 완료된 Task만 완료 이벤트 발행 단계로 전달 |
@MessagingGateway | 일반 Java 메서드 호출을 메시지 발행으로 변환 |
RequestHandlerRetryAdvice | Handler 재시도 및 최종 실패 복구 |
errorChannel | 분류되지 않은 메시징 오류의 최후 안전망 |
StreamBridge | SI outbound 메시지를 Cloud Stream 출력 바인딩으로 전달 |
Kafka inbound에서 SI Flow까지¶
Kafka 메시지는 bootstrap-workflow의 Spring Cloud Stream Consumer가 받는다. Consumer는 업무 로직을 직접 실행하지 않고 SI Channel로 위임한다.
@Bean("jobRegistration")
public Consumer<Message<JobTicket>> jobRegistrationConsumer(
@Qualifier(WorkflowJobChannels.JOB_REGISTRATION_SI)
MessageChannel siChannel) {
return message -> siChannel.send(message);
}
처리 경로는 다음과 같다.
Kafka job-registration 토픽
→ Spring Cloud Stream jobRegistration-in-0
→ jobRegistration Consumer Bean
→ jobRegistration.si DirectChannel
→ jobRegistrationFlow
→ Coordinator.startJob()
Consumer Bean과 Kafka 토픽 연결은 application-base.yml의 spring.cloud.function.definition과 spring.cloud.stream.bindings가 담당한다.
DirectChannel과 실행 스레드¶
워크플로우의 SI Channel은 대부분 DirectChannel이다.
@Bean(WorkflowJobChannels.JOB_EXECUTION_SI)
public MessageChannel jobExecutionSiChannel() {
return new DirectChannel();
}
DirectChannel은 메시지를 별도 스레드로 보내지 않는다. send()를 호출한 스레드에서 Subscriber와 IntegrationFlow가 동기적으로 실행된다.
Kafka Consumer 스레드
→ siChannel.send(message)
→ Router
→ Handler
→ Gateway / StreamBridge 발행 호출
→ Consumer 호출 복귀
따라서 Handler 예외는 같은 호출 스택을 통해 Consumer까지 전파될 수 있다. 새로운 비동기 처리 사이클은 메시지가 Kafka에 발행되고 다른 Consumer가 이를 수신할 때 시작된다.
Job IntegrationFlow¶
Job Flow는 WorkflowJobFlowConfig에 정의된다.
flowchart TD
A[jobRegistration.si] --> B[Coordinator.startJob]
C[jobExecution.si] --> D{JobStatus Router}
D -->|STARTED| E[job.onNext]
D -->|PAUSING| F[job.onPause]
E --> G[JobExecutor.onNext]
F --> H[JobExecutor.onPause]
I[jobCompletion.si] --> J{COMPLETED}
J --> K[JobExecutor.onComplete]
L[jobError.si] --> M{FAILED}
M --> N[JobExecutor.onError]
Job 등록 Flow는 Handler에서 Coordinator.startJob()을 호출한다.
IntegrationFlow.from(WorkflowJobChannels.JOB_REGISTRATION_SI)
.<JobTicket>handle((ticket, headers) -> {
coordinator.startJob(ticket);
return null;
})
.get();
Handler가 null을 반환하면 현재 Flow는 종료된다. 다음 라이프사이클은 startJob()이 발행한 jobExecution Kafka 이벤트를 새 Consumer가 수신하면서 시작된다.
Job 실행 Flow는 Job 상태를 내부 Channel로 라우팅한다.
IntegrationFlow.from(WorkflowJobChannels.JOB_EXECUTION_SI)
.<JobTicket, JobStatus>route(
ticket -> JobStatus.valueOf(ticket.getStatus()),
mapping -> mapping
.channelMapping(JobStatus.STARTED, WorkflowJobChannels.JOB_ON_NEXT)
.channelMapping(JobStatus.PAUSING, WorkflowJobChannels.JOB_ON_PAUSE))
.get();
Router는 업무를 수행하지 않고 목적 Channel만 선택한다. job.onNext, job.onPause에서 시작하는 별도 Flow가 각각 JobExecutor.onNext()와 onPause()를 호출한다.
Task IntegrationFlow¶
Task Flow는 WorkflowTaskFlowConfig에 정의된다.
flowchart TD
A[taskRegistration.si] --> B[Worker.startTask]
B --> C[taskExecution Kafka 발행]
D[taskExecution.si] --> E{executionBucket Router}
E -->|RUN| F[task.execution.run]
E -->|PAUSE| G[task.execution.pause]
E -->|UNSUPPORTED| H[task.execution.unsupported]
F --> I[TaskTimeoutRunner]
I --> J[TaskExecutor.onNext]
J --> K{COMPLETED Filter}
K -->|통과| L[TaskCompletion 발행]
K -->|탈락| M[nullChannel]
N[taskCompletion.si] --> O[Worker.completeTask]
O --> P{hasMoreTasks Router}
P -->|true| Q[다음 Task 등록]
P -->|false| R[JobCompletion 발행]
Task 실행 Router는 상태를 세 실행 Bucket으로 변환한다.
static String executionBucket(TaskTicket ticket) {
TaskStatus status = TaskStatus.valueOf(ticket.getStatus());
return switch (status) {
case STARTED, RESUMING -> "RUN";
case PAUSING -> "PAUSE";
default -> "UNSUPPORTED";
};
}
| Task 상태 | SI 실행 Channel |
|---|---|
STARTED, RESUMING | TASK_EXECUTION_RUN |
PAUSING | TASK_EXECUTION_PAUSE |
| 그 외 상태 | TASK_EXECUTION_UNSUPPORTED |
RUN Flow는 Timeout, Handler 실행, 완료 Filter, 완료 이벤트 발행을 연결한다.
IntegrationFlow.from(WorkflowTaskChannels.TASK_EXECUTION_RUN)
.<TaskTicket>handle((ticket, headers) -> {
taskTimeoutRunner.runWithTimeout(ticket, () -> taskExecutor.onNext(ticket));
return ticket;
}, endpoint -> endpoint.advice(taskRetryAdvice))
.filter(TaskTicket.class,
ticket -> TaskStatus.COMPLETED.label().equals(ticket.getStatus()),
filter -> filter.discardChannel("nullChannel"))
.<TaskTicket>handle((ticket, headers) -> {
gateway.sendToTasksCompletion(ticket);
return null;
}, endpoint -> endpoint.advice(taskRecovererAdvice))
.get();
TaskExecutor 호출이 정상 반환해도 Task가 반드시 COMPLETED인 것은 아니다. PAUSING이나 DELEGATING 상태라면 완료 이벤트를 발행하면 안 되므로 Filter가 COMPLETED만 통과시킨다. 탈락한 메시지는 아무 작업도 하지 않는 Spring Integration 기본 nullChannel로 보낸다.
Task 완료 Flow에서는 coordinator.hasMoreTasks()를 Router 조건으로 사용한다.
true
→ 다음 TaskTicket 생성
→ Coordinator.registerTask()
→ taskRegistration Kafka 이벤트
false
→ TaskTicket에서 JobTicket 생성
→ jobCompletion Kafka 이벤트
ControlTask IntegrationFlow¶
switch가 선택한 하위 Task는 WorkflowControlTaskFlowConfig의 별도 Flow를 사용한다. 구조는 일반 Task와 비슷하지만 ControlTask 전용 Executor, Channel, 완료·오류 토픽을 사용한다.
controlTaskRegistration.si
→ Worker.startControlTask()
→ controlTaskExecution
→ 상태 Router
→ TaskTimeoutRunner
→ ControlTaskExecutor.onNext()
→ controlTaskCompletion
→ hasMoreControlTasks Router
→ 다음 ControlTask 또는 부모 Task 완료 흐름
일반 Task와 ControlTask를 분리하면 재시도·오류 처리·메트릭을 계열별로 적용하고, ControlTask 체인 완료 후 부모 Task로 복귀하는 경로를 명확하게 관리할 수 있다.
WorkflowMessageGateway와 outbound Flow¶
업무 객체는 StreamBridge나 Kafka API를 직접 호출하지 않는다. WorkflowMessageGateway의 일반 Java 메서드를 호출한다.
@MessagingGateway
public interface WorkflowMessageGateway {
@Gateway(requestChannel = "taskExecution.out")
void sendToTasksExecution(TaskTicket taskTicket);
@Gateway(requestChannel = "taskCompletion.out")
void sendToTasksCompletion(TaskTicket taskTicket);
}
Spring Integration은 Gateway 호출을 Message 생성과 Channel 전송으로 변환한다.
gateway.sendToTasksCompletion(ticket)
→ Message<TaskTicket> 생성
→ taskCompletion.out Channel
→ WorkflowIntegrationConfig outbound Flow
→ StreamBridge.send(taskCompletion-out-0)
→ Spring Cloud Stream Binding
→ Kafka task-completion 토픽
이 구조는 Coordinator·Worker·Executor가 Kafka Binding 이름이나 메시징 API를 몰라도 되게 한다.
Retry Advice와 failure Channel¶
Task 실행 Handler에는 RequestHandlerRetryAdvice가 적용된다. 현재 기본 실행 정책은 최대 1회 재시도와 1초 지연이며, 다음 예외는 재시도 대상에서 제외된다.
WorkflowTimeoutExceptionWorkflowInterruptedExceptionWorkflowNonRetryableException
Handler 예외
→ RequestHandlerRetryAdvice
→ 재시도 가능하면 1회 재실행
→ 최종 실패
→ ErrorMessageSendingRecoverer
→ task.failure.si 또는 controlTask.failure.si
→ ErrorTicketFinalizer
→ TaskTicket FAILED 처리
→ taskError 또는 controlTaskError 발행
완료 이벤트 발행이나 PAUSING처럼 재시도가 부적절한 경로에는 재시도 횟수 0인 Recoverer 전용 Advice를 적용한다.
errorChannel 안전망¶
정상적인 Task 오류는 전용 task.failure.si 또는 controlTask.failure.si로 이동해야 한다. Advice 누락이나 잘못된 배선으로 예외가 Spring Integration 기본 errorChannel에 들어오는 경우를 대비해 WorkflowErrorFlowConfig가 최후 안전망을 제공한다.
분류되지 않은 메시징 오류
→ errorChannel
→ 원본 TaskTicket 추출
→ unclassified 오류 메트릭과 로그 기록
→ TaskTicket 실패 확정
→ taskError 발행
오류 토픽 발행 자체가 실패했을 때 다시 같은 errorChannel 경로로 무한 진입하지 않도록 발행 실패는 메트릭과 운영 복구 대상으로 종결한다.
동기·비동기 경계¶
sequenceDiagram
participant K as Kafka Consumer
participant C as DirectChannel
participant F as IntegrationFlow
participant H as Handler
participant G as Gateway/StreamBridge
participant T as Kafka Topic
K->>C: siChannel.send(message)
C->>F: 동기 전달
F->>H: Router/Handler 실행
H-->>F: 결과 또는 예외
F->>G: outbound 메시지 발행
G->>T: Kafka 전송
Note over T: 다음 Consumer 수신부터 새 비동기 처리 사이클
하나의 Job 전체가 단일 DB 트랜잭션으로 처리되는 것이 아니다. 각 Kafka 라이프사이클 이벤트가 독립적인 처리 및 장애 복구 단위가 된다.
Spring Integration 로그 해석¶
현재 로그는 Kafka와 SI 경계를 구분한다.
| 관찰 결과 | 의심 구간 |
|---|---|
[KAFKA 📨]가 없음 | Kafka 토픽, Binding, Consumer group |
[KAFKA 📨]는 있고 [SI ⚙️]가 없음 | Consumer에서 SI Channel로 가는 Bridge |
[SI ⚙️] 이후 완료가 없음 | Router, Handler, Timeout, TaskExecutor |
| 완료 발행은 있으나 다음 수신이 없음 | Outbound Flow, StreamBridge, Kafka 발행 |
[SI ⚠️] unclassified | Advice 누락 또는 errorChannel 오배선 |
Spring Integration 확장 기준¶
단순 업무 Task를 추가할 때는 IntegrationFlow를 수정하지 않는다. BaseTaskHandler 구현체를 Spring Bean으로 등록하고 파이프라인에서 해당 type을 사용한다.
새 Task 상태나 라이프사이클 이벤트를 추가할 때는 다음 연결 전체를 검토한다.
상태/이벤트 정의
→ Channel 상수와 MessageChannel Bean
→ Router 매핑
→ 상태별 IntegrationFlow
→ WorkflowMessageGateway
→ Outbound Flow와 StreamBridge
→ Cloud Stream Binding과 Kafka 토픽
→ Consumer Bean과 inbound SI Channel
→ Retry/errorChannel 정책
→ Flow 통합 테스트
Spring Integration에는 메시지 전달, 상태 라우팅, Handler 연결, 완료 필터링, 재시도·복구와 관측 로직을 둔다. 실제 업무 규칙은 TaskHandler, Coordinator, Worker 같은 도메인 객체에 유지한다.
Task 실행 과정¶
Task는 다음 순서로 처리된다.
TaskRegistration
→ Task DB 등록
→ TaskExecution
→ pre Task 실행
→ type에 해당하는 TaskHandler 실행
→ 반환값을 Context에 저장
→ post Task 실행
→ 상태 COMPLETED
→ finalize Task 실행
→ TaskCompletion
Task 완료 후 Coordinator.hasMoreTasks()가 다음 Task 존재 여부를 판단한다. 다음 Task가 있으면 등록하고, 더 없으면 JobCompletion을 발행한다. 예외가 발생하면 TaskError를 거쳐 Job 실패 흐름으로 연결된다.
Task type과 Spring Bean 이름¶
파이프라인의 type은 Spring Bean 이름이다.
DefaultTaskHandlerResolver는 ApplicationContext.getBean(taskType)으로 Handler를 찾으므로 두 문자열이 반드시 같아야 한다. 등록되지 않은 타입을 사용하면 Bean을 찾지 못해 Task가 실패한다.
현재 기본 예시는 다음과 같다.
| type | 기능 |
|---|---|
io/print | 로그 출력 |
io/print-resume | 일시정지·재개 예제 |
random/int | 임의 정수 생성 |
time/sleep | 일정 시간 대기 |
switch | 조건 분기 |
신규 TaskHandler 만들기¶
문자열을 대문자로 변환하는 Task 예제다.
package kr.co.dandisoft.minu.workflow.task.handler.text;
import kr.co.dandisoft.minu.workflow.task.handler.BaseTaskHandler;
import kr.co.dandisoft.minu.workflow.task.handler.TaskHandler;
import kr.co.dandisoft.minu.workflow.ticket.TaskTicket;
import org.springframework.stereotype.Component;
@Component("text/uppercase")
public class Uppercase extends BaseTaskHandler implements TaskHandler<Object> {
@Override
protected Object doHandle(TaskTicket taskTicket) {
String value = taskTicket.getRequiredString("value");
return value.toUpperCase();
}
}
파이프라인에서는 다음과 같이 사용한다.
tasks:
- type: text/uppercase
label: 이름 대문자 변환
value: "${yourName}"
output: upperName
- type: io/print
label: 변환 결과 출력
text: "변환 결과: ${upperName}"
단순 Handler는 Resolver를 수정할 필요 없이 @Component("type 문자열")로 등록하면 된다.
switch와 ControlTask¶
tasks:
- type: switch
label: 실행 경로 선택
expression: "${selector}"
cases:
- key: hello
tasks:
- type: io/print
text: "Hello path"
- key: bye
tasks:
- type: io/print
text: "Goodbye path"
default:
- tasks:
- type: io/print
text: "Unknown path"
switch는 표현식을 평가한 후 선택된 Task를 ControlTask 흐름으로 위임한다.
SwitchTaskHandler
→ status=DELEGATING
→ ControlTaskRegistration
→ ControlTaskExecution
→ 선택된 Task 실행
→ 원래 Job 흐름으로 복귀
처음에는 일반 순차 Task로 흐름을 익힌 후 분기 기능을 추가하는 편이 안전하다.
4. 핵심 패턴 & 안티패턴¶
권장 패턴¶
- Task의 반환값은
output이름으로 Context에 저장하고 다음 Task에서${변수명}으로 참조한다. - Task
type과@ComponentBean 이름을 동일하게 유지한다. - 업무 로직은 TaskHandler에 두고, Job·Task 순서 제어는 Coordinator와 이벤트 흐름에 맡긴다.
- Handler는 같은 Kafka 메시지가 다시 전달될 수 있음을 고려해 가능한 한 멱등하게 작성한다.
- 처음에는 순차 Task로 검증하고 필요한 경우에만
switch와 ControlTask를 추가한다.
피해야 할 패턴¶
- 등록되지 않은 임의의
type을 파이프라인에 사용하지 않는다. - TaskHandler가 다음 Task를 직접 호출하거나 Kafka 토픽을 임의로 발행하지 않는다.
- 공유해야 할 실행 값을 전역 변수나 싱글턴 Bean 필드에 저장하지 않는다.
- 비동기 실행 API의
200 OK를 워크플로우 완료로 판단하지 않는다. - 외부 시스템 호출 Handler를 재시도와 멱등성 고려 없이 구현하지 않는다.
5. 트러블슈팅¶
로컬 실행과 확인¶
필요한 주요 환경은 Java 25, PostgreSQL, Kafka이며 프로파일에 따라 Redis, Vault, Eureka가 추가로 필요할 수 있다.
기본 서비스 포트는 18020, 로그 경로는 _temp/logs/다.
./gradlew :hera-workflow:test
./gradlew :bootstrap-workflow:test
./gradlew :bootstrap-workflow:build
실행 시 로그에서 다음 이벤트를 확인한다.
주요 조회 API는 다음과 같다.
GET /api/v1/wfw/job/{jobId}
GET /api/v1/wfw/task/{taskId}
GET /api/v1/wfw/history/jobs/pageable
GET /api/v1/wfw/history/task/pageable
Kafka가 정상 동작하지 않으면 실행 API가 성공했더라도 후속 이벤트가 진행되지 않을 수 있으므로 애플리케이션 로그, Kafka 연결 상태, wfw_job과 wfw_task 실행 이력을 함께 확인한다.
자주 확인할 문제¶
| 증상 | 확인 사항 |
|---|---|
| 실행 API는 성공했지만 Task가 시작되지 않음 | Kafka 연결, JobRegistration/JobExecution 토픽, Consumer group 로그 |
NoSuchBeanDefinitionException 발생 | 파이프라인 type과 @Component 이름 일치 여부 |
${변수}가 치환되지 않음 | 앞 Task의 output, 입력 파라미터 이름, Context 저장 여부 |
| Task가 FAILED로 종료됨 | _temp/logs/, wfw_task, 오류 로그와 TaskError 흐름 |
| 스케줄 워크플로우가 실행되지 않음 | trigger, scheduler, activeYn 및 스케줄 재등록 로그 |
| API 호출 직후 결과가 없음 | 정상적인 비동기 처리인지 Job/Task 이력 API로 확인 |
권장 실습 순서¶
io/print하나만 있는 파이프라인 등록/api/v1/wfw/run으로 실행- Job 및 Task 이력 확인
input과${변수}사용- 반환값을
output으로 Context에 저장 pre,post,finalize추가- 사용자 정의 TaskHandler 구현
- 예외를 발생시켜 Task/Job 실패 흐름 확인
switch와 ControlTask 실습- CRON 또는 고정 주기 스케줄 추가
6. 레퍼런스 링크¶
주요 소스 위치¶
| 역할 | 파일 |
|---|---|
| 실행 API | hera-workflow/.../web/api/wfw/WorkflowApiController.java |
| 설정 등록 API | hera-workflow/.../web/api/wfw/WorkflowConfigApiController.java |
| Job/Task 조정 | hera-workflow/.../workflow/Coordinator.java |
| Task 상태 처리 | hera-workflow/.../workflow/Worker.java |
| Job 이벤트 흐름 | hera-workflow/.../integration/WorkflowJobFlowConfig.java |
| Task 이벤트 흐름 | hera-workflow/.../integration/WorkflowTaskFlowConfig.java |
| Task 실행기 | hera-workflow/.../task/executor/DefaultTaskExecutor.java |
| Handler 검색 | hera-workflow/.../task/resolver/DefaultTaskHandlerResolver.java |
| Kafka 바인딩 | bootstrap-workflow/src/main/resources/application-base.yml |
docs/references/WORKFLOW.md는 제목과 일부 구현이 HERA-V3 기준이므로 개념 참고로만 사용하고, 실제 동작은 현재 HERA-V4 소스와 설정을 우선한다.