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참조.
댓글 (0)
로그인하면 이 기사에 내 생각을 남길 수 있어요