콘텐츠로 이동

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에 저장된다.

POST /api/v1/wfw/config
Content-Type: application/json

주요 요청 필드는 다음과 같다.

{
  "workflowId": "sample-workflow",
  "workflowName": "샘플 워크플로우",
  "trigger": "1",
  "scheduler": "",
  "activeYn": "Y",
  "description": "워크플로우 예제",
  "pipeline": "파이프라인 YAML 문자열"
}

workflowId를 생략하면 내부에서 ID가 생성된다. 등록 시 ConfigBuilderpipeline 문자열을 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 실행 전 수행할 보조 Task
  • tasks: 실제 업무 단계
  • post: 본 Task가 정상 처리된 후 수행
  • finalize: Task 실행 마무리 단계
  • outputs: Job 종료 시 노출할 결과
  • ${변수명}: Context 값 참조
  • output: Task 반환값을 Context에 저장할 이름

예를 들어 random/int42를 반환하고 output: randomNumber가 지정되어 있으면 Context에 randomNumber=42가 저장된다. 이후 Task에서 ${randomNumber}로 참조한다.

워크플로우 실행

POST /api/v1/wfw/run
Content-Type: application/x-www-form-urlencoded
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.ymlspring.cloud.function.definitionspring.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초 지연이며, 다음 예외는 재시도 대상에서 제외된다.

  • WorkflowTimeoutException
  • WorkflowInterruptedException
  • WorkflowNonRetryableException
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 Consumer가 메시지를 수신
[SI ⚙️] Spring Integration Flow가 메시지를 처리
[SI ⚠️] 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 이름이다.

- type: io/print
  text: "Hello"
@Component("io/print")
public class Print extends BaseTaskHandler {
}

DefaultTaskHandlerResolverApplicationContext.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@Component Bean 이름을 동일하게 유지한다.
  • 업무 로직은 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가 추가로 필요할 수 있다.

./gradlew :bootstrap-workflow:bootRun

기본 서비스 포트는 18020, 로그 경로는 _temp/logs/다.

./gradlew :hera-workflow:test
./gradlew :bootstrap-workflow:test
./gradlew :bootstrap-workflow:build

실행 시 로그에서 다음 이벤트를 확인한다.

[REGISTER_JOB]
[STARTED_JOB]
[STARTED_TASK]
[COMPLETED_TASK]

주요 조회 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_jobwfw_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로 확인

권장 실습 순서

  1. io/print 하나만 있는 파이프라인 등록
  2. /api/v1/wfw/run으로 실행
  3. Job 및 Task 이력 확인
  4. input${변수} 사용
  5. 반환값을 output으로 Context에 저장
  6. pre, post, finalize 추가
  7. 사용자 정의 TaskHandler 구현
  8. 예외를 발생시켜 Task/Job 실패 흐름 확인
  9. switch와 ControlTask 실습
  10. 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 소스와 설정을 우선한다.