TL;DR

  • AdaptiveFlow는 Java 21 가상 스레드에서 작업 방향성 비순환 그래프(DAG)를 실행하는 런타임 의존성 없는 소규모 라이브러리로, 워크플로를 검증·정렬하고 결과를 불변 타입 안전 객체로 반환함.
  • 플루언트 API로 작업을 한 번에 정의하며, 병렬 그룹과 의존 관계를 바탕으로 실행 계획을 구성함.
  • DAG 계획기는 중복 ID, 알 수 없는 의존성, 사이클을 검사하고, 실행 엔진은 작업별 CompletableFuture 게이트와 가상 스레드로 작업을 실행함.
  • 재시도는 고정 지연 또는 지수 백오프를 지원하며, 지수 백오프에는 전체 지터를 적용하고 재시도 소진 시 다운스트림 작업을 건너뜀.
  • Java 21에서 선형 5개 작업 워크플로는 초당 약 5,300회, 팬인 DAG 워크플로는 초당 약 6,500회를 기록함.

AdaptiveFlow

  • AdaptiveFlow는 가상 스레드에서 작업 DAG를 실행하는 Java 21용 런타임 의존성 없는 라이브러리임.
  • 플루언트 API로 워크플로를 한 번에 구성하면 엔진이 검증과 위상 정렬을 수행하고, 작업마다 가상 스레드 하나를 사용해 실행한 뒤 불변의 타입 안전 결과 객체를 반환함.

배경

  • JVM에서 워크플로 오케스트레이션을 구현할 때는 대개 무거운 프레임워크가 필요하지만, AdaptiveFlow는 그 반대 방향의 라이브러리임.
  • 공개 타입을 15개 미만으로 제한하고, Java 21의 가상 스레드, 레코드, CompletableFuture를 기반으로 함.
  • Spring, XML, 에이전트, 관리해야 할 스레드 풀이 필요하지 않으며, 캐리어 스레드는 JVM이 처리함.

빠른 시작

설치

  • Maven 좌표: io.github.varun-51:adaptiveflow-core:1.0.2
  • Gradle 의존성: io.github.varun-51:adaptiveflow-core:1.0.2

사용

  • WorkflowBuilder.builder("etl")로 워크플로를 만들고, extract, validate, summary, archives, load 작업을 정의하는 방식임.
  • 작업은 앞서 추가된 모든 작업에 의존하며, parallel(...) 그룹에 넣은 TaskRef들은 동시에 실행됨. 그룹 다음 작업은 그룹 안의 모든 작업에 의존함.
  • RetryPolicy.exponentialBackoff(4, Duration.ofMillis(100))처럼 재시도 정책을 설정하고 execute()를 호출하면 DAG 구성, 검증, 위상 정렬, 실행을 거쳐 불변 WorkflowResult를 반환함.
  • 결과가 성공이면 result.<Summary>result("summary")처럼 타입을 지정해 작업 결과를 조회함.
  • 공개 API는 이 전체 흐름으로 구성되며, 추가로 학습해야 할 요소가 없음.

작동 방식

  • WorkflowBuilder가 TaskSpec 목록을 만들고, DagPlanner가 이를 검증하고 위상 정렬해 ExecutionPlan으로 고정함.
  • ExecutionEngine은 작업마다 게이트를 두고, CompletableFuture 팬인을 거쳐 가상 스레드에서 작업을 실행하며, 결과를 WorkflowResult에 기록함.
  • DagPlanner는 중복 ID, 알 수 없는 의존성, Kahn 알고리즘으로 판별하는 사이클을 검증하고 불변 ExecutionPlan을 만듦.
  • 다중 루트 DAG도 허용되며, 이를 통해 병렬 진입점을 표현함.
  • ExecutionEngine은 모든 작업에 CompletableFuture 게이트를 할당함. 모든 의존성 게이트가 완료된 뒤 작업을 실행기에 제출함.
  • 작업마다 가상 스레드 하나를 사용하므로 의존 작업과 독립 분기가 머신 자원에 자연스럽게 맞춰짐.
  • 재시도는 해당 작업의 가상 스레드 안에서 수행됨. 실패한 작업은 성공하거나 시도 횟수를 모두 소진할 때까지 설정된 고정 또는 지수 백오프로 다시 호출됨.
  • 지수 백오프 대기에는 전체 지터를 적용해 동기화된 실패가 같은 시한에 겹치지 않도록 함.
  • 작업이 재시도를 모두 소진하면 실행은 다운스트림 작업의 디스패치를 중단함. 아직 대기 중인 작업은 모두 SKIPPED로 완료되며 taskResult(id)로 조회 가능함.
  • 이 경우 WorkflowResult.isSuccess()는 false이고, 실패한 작업의 TaskResult에서 오류를 확인할 수 있음. 실행 중인 형제 작업은 중단되지 않고 정상적으로 완료되어 결과를 보고함.
  • 결과는 작업 ID를 키로 사용하는 스레드 안전 ExecutionContext에 저장되며, ctx.<T>result(id)는 호출 지점에서 타입을 확인함.
  • 실행이 끝나면 WorkflowResult가 불변 스냅샷을 취함.

재시도

  • none(): 재시도 없음(기본값).
  • fixedDelay(attempts, delay): 시도 사이에 고정된 시간 동안 대기함.
  • exponentialBackoff(attempts, firstDelay, maxDelay): 시도마다 대기 시간이 두 배가 되며 maxDelay를 상한으로 둠.
  • 시도 횟수를 모두 사용하기 전의 실패는 호출자에게 노출되지 않으며, 재시도 소진 시 마지막 오류만 노출됨.

실행기

  • 작업별 가상 스레드가 기본 실행 방식임.
  • 사용자 실행기를 지정할 수 있으며, 예를 들어 workflow.execute(Executors.newFixedThreadPool(4))를 사용할 수 있음.
  • 라이브러리는 호출자가 소유한 실행기를 종료하지 않음. 인자 없는 execute()는 실행기를 소유하고 종료함.

벤치마크

  • JMH 제품군은 adaptiveflow-core/src/test/java/io/github/varun51/adaptiveflow/bench에서 실행됨.
  • -f 0은 셰이딩된 벤치마크 JAR 없이 실행할 수 있도록 인프로세스 모드를 사용함. CI 수준의 격리가 필요하면 포크 수를 5 이상으로 설정함.
  • 측정 환경은 Java 21, 인프로세스 모드, 1초씩 3회 반복임.
  • 선형 작업 5개 워크플로: 초당 약 5,300회.
  • 팬인 DAG 작업 10개 워크플로: 초당 약 6,500회.
  • 측정 대상은 워크플로 구성, 검증, 계획, 실행, 결과 수집까지의 종단 간 작업이며 가상 스레드도 포함함.

소스에서 빌드

  • ./mvnw -pl adaptiveflow-core verify로 전체 품질 검증을 실행함.
  • 검증 항목은 동시성·재시도·검증 관련 단위 테스트 61개, JaCoCo 라인·분기 커버리지 기준 90% / 80%, 사용자 정의 Checkstyle 규칙 세트, 최대 노력 수준의 SpotBugs임.

요구 사항

  • Java 21이 필요하며, 가상 스레드는 정식 기능이므로 프리뷰 플래그가 필요하지 않음.

아직 지원하지 않는 기능

  • 작업 시간 제한 및 워크플로 전체 시간 제한
  • 실행 중인 워크플로의 취소 및 인터럽트
  • 컨텍스트의 사용자 정의 변수(작업 출력만 저장)
  • 내구성 실행 및 체크포인팅(이 문제 영역은 Temporal 참조)

로드맵

  • v1.0은 결과의 결정성을 유지하며, 같은 워크플로와 입력은 언제나 같은 결과를 생성함.
  • 무작위성이 안전한 부분에는 이미 적응형 처리가 적용되어 있음. 지수 재시도 대기에 전체 지터를 적용해 동기화된 실패가 같은 백오프 시한에 몰리지 않도록 함.
  • 자원 인식 및 실시간 실행 학습을 포함한 적응형 스케줄링은 v1.0 이후 방향임.
  • 플러그형 실행기와 작업별 재시도 정책이 그 방향을 확장하는 접점임.

라이선스

  • Apache License 2.0이며, 자세한 내용은 LICENSE 참조.