학습 목표
@ActivityInterface/@ActivityMethod로 액티비티를 정의하고, 워크플로우와 달리 아무 제약이 없다는 것을 확인한다ActivityOptions를 워크플로우 코드 안에서 만들고Workflow.newActivityStub으로 스텁을 얻는다- 타임아웃 4종(ScheduleToStart / StartToClose / ScheduleToClose / Heartbeat)이 각각 어느 구간을 재고 어떤 장애를 잡는지 타임라인으로 구분하고, 터졌을 때의 히스토리 이벤트를 JSON 으로 확인한다
- Heartbeat 로 진행 상황을 저장해, 10만 건 배치를 5만 건에서 죽였다가 5만 번째부터 재개한다
- Temporal 의 at-least-once 보장 때문에 결제가 두 번 되는 시나리오를 재현하고 멱등키로 막는다
- Local Activity 와 일반 Activity 의 히스토리 비용 차이를 실측한다
선행 스텝: Step 03 — 워크플로우 정의와 결정성 예상 소요: 120분
docker compose ps
temporal workflow list --limit 3결과
NAME IMAGE STATUS PORTS
temporal temporalio/auto-setup:1.22.4 Up 12 minutes 0.0.0.0:7233->7233/tcp
temporal-ui temporalio/ui:2.21.3 Up 12 minutes 0.0.0.0:8233->8080/tcp
temporal-postgres postgres:15 Up 12 minutes 5432/tcp
Status WorkflowId Type StartTime
COMPLETED order-1001 OrderWorkflow 2026-03-11T09:12:04Z
COMPLETED order-1003 OrderWorkflow 2026-03-11T09:41:02Z
RUNNING order-2002 OrderWorkflow 2026-03-11T19:04:11Z@ActivityInterface — 제약이 없는 세계Step 03 에서 워크플로우 코드가 할 수 없는 일의 목록을 봤습니다. 액티비티는 그 정반대입니다. 평범한 Java 코드입니다.
package com.example.order;
import io.temporal.activity.ActivityInterface;
import io.temporal.activity.ActivityMethod;
@ActivityInterface
public interface PaymentActivity {
@ActivityMethod
String charge(String orderId, long amount);
@ActivityMethod
void refund(String paymentId);
}public class PaymentActivityImpl implements PaymentActivity {
private static final Logger log = LoggerFactory.getLogger(PaymentActivityImpl.class);
private final HttpClient http = HttpClient.newHttpClient(); // 필드에 커넥션 풀 ─ 괜찮습니다
@Override
public String charge(String orderId, long amount) {
String requestId = UUID.randomUUID().toString(); // 랜덤 ─ 괜찮습니다
long start = System.currentTimeMillis(); // 벽시계 ─ 괜찮습니다
HttpRequest req = HttpRequest.newBuilder()
.uri(URI.create("https://pg.example.com/v1/charges"))
.header("Idempotency-Key", requestId)
.POST(HttpRequest.BodyPublishers.ofString(body(orderId, amount)))
.build();
try {
HttpResponse<String> res = http.send(req, HttpResponse.BodyHandlers.ofString()); // 네트워크 ─ 괜찮습니다
log.info("PG 응답 {}ms status={}", System.currentTimeMillis() - start, res.statusCode());
return parsePaymentId(res.body());
} catch (Exception e) {
throw ApplicationFailure.newFailure("PG 호출 실패", "PaymentGatewayError", e.getMessage());
}
}
}| 항목 | 워크플로우 | 액티비티 |
|---|---|---|
@Xxx Interface 애너테이션 | @WorkflowInterface | @ActivityInterface |
| 메서드 개수 | @WorkflowMethod 정확히 하나 | @ActivityMethod 여러 개 가능 |
| 등록 방식 | 클래스를 등록 (registerWorkflowImplementationTypes) | 인스턴스를 등록 (registerActivitiesImplementations) |
| 인스턴스 생명주기 | 실행마다 새로 생성 | Worker 당 하나를 공유 |
| 벽시계·랜덤·스레드 | 금지 | 허용 |
| 네트워크·파일·DB | 금지 | 허용 |
| Spring 의존성 주입 | 사실상 금지 | 권장 |
| 재실행 | 리플레이로 몇 번이든 | 실패 시 재시도 (기본 무한) |
@ActivityMethod 는 사실 생략 가능합니다. @ActivityInterface 가 붙은 인터페이스의 모든 public 메서드가 자동으로 액티비티가 됩니다. 이름을 바꾸고 싶을 때만 붙입니다.
@ActivityInterface
public interface InventoryActivity {
@ActivityMethod(name = "ReserveStock") // 기본값은 "Reserve" (메서드명 첫 글자 대문자)
String reserve(String orderId, String sku, int qty);
void release(String reservationId); // 애너테이션 없어도 액티비티. 타입명은 "Release"
}⚠️ 함정 — 액티비티 메서드명을 바꾸면 실행 중인 워크플로우가 깨집니다 액티비티 타입 이름은 히스토리에 문자열로 저장됩니다.
"activityType": { "name": "Charge" }
charge를chargeWithFee로 리팩터링하고 배포하면, 히스토리에Charge가 적힌 워크플로우를 리플레이할 때 SDK 가 그 타입을 찾지 못합니다.io.temporal.failure.ApplicationFailure: Activity Type "Charge" is not registered with a worker. Known types are: ChargeWithFee, Refund, Reserve, Release해결은
@ActivityMethod(name = "Charge")로 옛 이름을 못 박는 것입니다. 자바 메서드명은 자유롭게 바꾸되 액티비티 타입명은 고정하세요. 본격적인 대응은 Step 10(버저닝)에서 다룹니다.
ActivityOptions 와 스텁 생성워크플로우는 액티비티를 직접 호출하지 않습니다. 스텁(stub) 을 통해 호출하고, 스텁은 실제로는 "이 액티비티를 스케줄해 달라"는 Command 를 만듭니다.
public class OrderWorkflowImpl implements OrderWorkflow {
// ActivityOptions 는 워크플로우 코드 안에서 만든다.
// Worker 설정이 아니라 "이 워크플로우가 이 액티비티를 어떻게 부를 것인가"의 선언이다.
private final PaymentActivity payment = Workflow.newActivityStub(
PaymentActivity.class,
ActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofSeconds(30))
.setScheduleToStartTimeout(Duration.ofMinutes(5))
.setRetryOptions(RetryOptions.newBuilder()
.setMaximumAttempts(5)
.build())
.build());
// ...
}옵션이 워크플로우 코드 안에 있다는 사실이 중요합니다. 같은 액티비티라도 워크플로우마다 다른 타임아웃을 줄 수 있고, 옵션을 바꾸면 그것도 히스토리에 기록됩니다.
temporal workflow show -w order-1001 --output json | jq '.events[4].activityTaskScheduledEventAttributes'결과
{
"activityId": "5",
"activityType": { "name": "Charge" },
"taskQueue": { "name": "ORDER_TASK_QUEUE", "kind": "TASK_QUEUE_KIND_NORMAL" },
"input": { "payloads": [ { "data": "IjEwMDEi" }, { "data": "MzkwMDA=" } ] },
"scheduleToCloseTimeout": "0s",
"scheduleToStartTimeout": "300s",
"startToCloseTimeout": "30s",
"heartbeatTimeout": "0s",
"workflowTaskCompletedEventId": "4",
"retryPolicy": {
"initialInterval": "1s",
"backoffCoefficient": 2,
"maximumInterval": "100s",
"maximumAttempts": 5
}
}"0s" 는 "설정하지 않음 = 무제한"이라는 뜻입니다.
액티비티별로 다른 옵션이 필요하면 메서드 단위로 지정합니다.
private final PaymentActivity payment = Workflow.newActivityStub(
PaymentActivity.class,
ActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofSeconds(30))
.build(),
Map.of(
// refund 는 더 느려도 되고, 더 오래 재시도해야 한다
"Refund", ActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofMinutes(2))
.setScheduleToCloseTimeout(Duration.ofHours(24))
.build()));이 절이 Step 04 의 하이라이트입니다. 액티비티 타임아웃은 네 가지이고, 재는 구간이 전부 다릅니다.
워크플로우가 액티비티를 호출
│
▼
t0 ─── ActivityTaskScheduled ─────────────────────────────────────────────────►
│
│ ◄─── ScheduleToStart ───►│
│ (Task Queue 에서 │
│ 대기하는 시간) │
│ │
│ t1 ─── ActivityTaskStarted ─────────────────────────►
│ │
│ │◄──────── StartToClose ────────►│
│ │ (한 번의 시도가 도는 시간) │
│ │ │
│ │ ♥ ♥ ♥ ♥ │
│ │ ◄─► ← Heartbeat 간격
│ │ t2 ─── Completed / TimedOut
│ │
│◄──────────────────── ScheduleToClose ──────────────────────►│
│ (재시도까지 전부 포함한 총 시간. 시도 1 + 대기 + 시도 2 + ... )
│
│ 재시도가 일어나면 이렇게 반복된다:
│ [대기][시도1][백오프][대기][시도2][백오프][대기][시도3] ...
│ ↑ ↑ ↑
│ S2S S2C ─────────────────────────────────────────►│
│ StartToClose 는 매 시도마다 새로 잰다| 타임아웃 | 재는 구간 | 기본값 | 잡아내는 장애 | 재시도 | 필수 여부 |
|---|---|---|---|---|---|
| ScheduleToStart | 큐에 들어간 순간 ~ Worker 가 집어간 순간 | 무한 | Worker 부족, Worker 전부 다운, 잘못된 Task Queue 이름 | 재시도 안 됨 (즉시 실패) | 선택 |
| StartToClose | 한 번의 시도가 시작 ~ 끝 | 무한 | 액티비티가 매달림, 외부 API 무응답, 무한 루프 | 재시도됨 | 이것 또는 S2C 필수 |
| ScheduleToClose | 첫 스케줄 ~ 최종 완료 (모든 재시도 포함) | 무한 | 전체가 너무 오래 걸림. 재시도 총량 제한 | (총량이므로 해당 없음) | 이것 또는 S2C 필수 |
| Heartbeat | 하트비트 사이의 간격 | 없음 | 조용히 죽은 Worker (프로세스 kill, OOM, 네트워크 단절) | 재시도됨 | 장기 액티비티에 강력 권장 |
둘 다 지정하지 않으면 워크플로우가 실행되는 순간 실패합니다.
// ✘ 타임아웃을 하나도 안 줬다
private final PaymentActivity payment = Workflow.newActivityStub(
PaymentActivity.class,
ActivityOptions.newBuilder().build());결과 (Worker 콘솔)
09:55:12.331 [workflow-method-order-3001-...] WARN i.t.i.s.WorkflowExecutionHandler - Workflow execution failure
io.temporal.failure.ApplicationFailure: message='Either ScheduleToCloseTimeout or StartToCloseTimeout is required', type='java.lang.IllegalStateException', nonRetryable=false
at io.temporal.internal.sync.ActivityStubBase.execute(ActivityStubBase.java:47)
at io.temporal.internal.sync.ActivityInvocationHandler.lambda$getActivityFunc$0(ActivityInvocationHandler.java:82)
at io.temporal.internal.sync.ActivityInvocationHandlerBase$ActivityStubInvocationHandler.invoke(ActivityInvocationHandlerBase.java:73)
at jdk.proxy2/jdk.proxy2.$Proxy14.charge(Unknown Source)
at com.example.order.OrderWorkflowImpl.processOrder(OrderWorkflowImpl.java:41)
Caused by: java.lang.IllegalStateException: Either ScheduleToCloseTimeout or StartToCloseTimeout is required
at io.temporal.common.interceptors.ActivityOptionsUtils.validateAndBuildOptions(ActivityOptionsUtils.java:61)
at io.temporal.internal.sync.SyncWorkflowContext.executeActivity(SyncWorkflowContext.java:265)
... 4 common frames omittedtemporal workflow describe -w order-3001결과
Execution Info:
Workflow Id order-3001
Type OrderWorkflow
Status FAILED
History Length 5
Failure:
Message Either ScheduleToCloseTimeout or StartToCloseTimeout is required
Type java.lang.IllegalStateException이 경우는 워크플로우 자체가 FAILED 입니다. Step 03 의 NonDeterministicException 과 달리 코드 로직 예외로 취급되기 때문입니다.
ActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofSeconds(3)) // 액티비티는 10초 걸린다
.setRetryOptions(RetryOptions.newBuilder().setMaximumAttempts(2).build())
.build();temporal workflow show -w order-3002결과
Progress:
ID Time Type
1 2026-03-11T10:02:11Z WorkflowExecutionStarted
2 2026-03-11T10:02:11Z WorkflowTaskScheduled
3 2026-03-11T10:02:11Z WorkflowTaskStarted
4 2026-03-11T10:02:11Z WorkflowTaskCompleted
5 2026-03-11T10:02:11Z ActivityTaskScheduled
6 2026-03-11T10:02:11Z ActivityTaskStarted
7 2026-03-11T10:02:19Z ActivityTaskTimedOut ← 3초 × 2회 시도 후
8 2026-03-11T10:02:19Z WorkflowTaskScheduled
9 2026-03-11T10:02:19Z WorkflowTaskStarted
10 2026-03-11T10:02:19Z WorkflowTaskCompleted
11 2026-03-11T10:02:19Z WorkflowExecutionFailed
Result:
Status: FAILED7번 이벤트의 상세를 봅니다.
temporal workflow show -w order-3002 --output json \
| jq '.events[] | select(.eventType=="EVENT_TYPE_ACTIVITY_TASK_TIMED_OUT")'결과
{
"eventId": "7",
"eventTime": "2026-03-11T10:02:19.114Z",
"eventType": "EVENT_TYPE_ACTIVITY_TASK_TIMED_OUT",
"activityTaskTimedOutEventAttributes": {
"failure": {
"message": "activity StartToClose timeout",
"timeoutFailureInfo": {
"timeoutType": "TIMEOUT_TYPE_START_TO_CLOSE"
}
},
"scheduledEventId": "5",
"startedEventId": "6",
"retryState": "RETRY_STATE_MAXIMUM_ATTEMPTS_REACHED"
}
}TIMEOUT_TYPE_START_TO_CLOSE 와 RETRY_STATE_MAXIMUM_ATTEMPTS_REACHED 두 값이 핵심입니다.
Worker 를 전부 내린 뒤 워크플로우를 시작합니다.
ActivityOptions.newBuilder()
.setScheduleToStartTimeout(Duration.ofSeconds(10))
.setStartToCloseTimeout(Duration.ofSeconds(30))
.build();결과
{
"eventId": "6",
"eventTime": "2026-03-11T10:11:32.508Z",
"eventType": "EVENT_TYPE_ACTIVITY_TASK_TIMED_OUT",
"activityTaskTimedOutEventAttributes": {
"failure": {
"message": "activity ScheduleToStart timeout",
"timeoutFailureInfo": {
"timeoutType": "TIMEOUT_TYPE_SCHEDULE_TO_START"
}
},
"scheduledEventId": "5",
"startedEventId": "0",
"retryState": "RETRY_STATE_NON_RETRYABLE_FAILURE"
}
}두 가지를 주목하세요.
"startedEventId": "0" — 시작조차 못 했습니다. ActivityTaskStarted 이벤트 자체가 없습니다."retryState": "RETRY_STATE_NON_RETRYABLE_FAILURE" — 재시도되지 않습니다. 재시도해도 어차피 같은 큐에 들어가 같은 이유로 대기할 것이기 때문입니다.{
"eventId": "14",
"eventType": "EVENT_TYPE_ACTIVITY_TASK_TIMED_OUT",
"activityTaskTimedOutEventAttributes": {
"failure": {
"message": "activity ScheduleToClose timeout",
"timeoutFailureInfo": {
"timeoutType": "TIMEOUT_TYPE_SCHEDULE_TO_CLOSE",
"lastHeartbeatDetails": null
},
"cause": {
"message": "activity StartToClose timeout",
"timeoutFailureInfo": { "timeoutType": "TIMEOUT_TYPE_START_TO_CLOSE" }
}
},
"scheduledEventId": "5",
"startedEventId": "13",
"retryState": "RETRY_STATE_TIMEOUT"
}
}cause 에 직전 시도가 왜 실패했는지가 중첩되어 들어갑니다. 세 번째 시도가 StartToClose 로 터지던 도중 전체 예산(ScheduleToClose)이 소진된 것입니다.
💡 실무 팁 — 어떻게 정하나
- StartToClose: 액티비티의 p99 응답시간 × 2~3배. 외부 API 라면 그쪽 SLA + 여유.
- ScheduleToStart: "이 정도 큐 대기는 이미 사고다" 싶은 값. 보통 1~5분. 알림용 지표로 쓰기 좋습니다.
- ScheduleToClose: 비즈니스 데드라인. "결제는 아무리 재시도해도 30분 안에 끝나야 한다."
- Heartbeat: 하트비트 주기의 2~3배. 30초마다 뛴다면 90초.
셋 다 주는 것이 원칙이지만, 하나만 고르라면 StartToClose 입니다.
⚠️ 함정 — StartToClose 만 넉넉히 주고 ScheduleToStart 를 안 주면
ActivityOptions.newBuilder() .setStartToCloseTimeout(Duration.ofMinutes(10)) // 넉넉하게! .build(); // ScheduleToStart 없음 = 무한배포 사고로 Worker 가 전부 죽었다고 합시다. 액티비티 태스크는 Task Queue 에 쌓입니다. StartToClose 는 시작한 뒤부터 재는 타이머이므로 아직 작동하지 않습니다. ScheduleToStart 가 무한이니 아무도 알려 주지 않습니다. 결과
$ temporal task-queue describe --task-queue ORDER_TASK_QUEUE --task-queue-type activity BuildId TaskQueueType Pollers BacklogCount ApproximateBacklogAge ACTIVITY 0 84291 3h 42m 11s워크플로우는 전부 Running, 액티비티는 3시간 42분째 큐에서 대기 중입니다. 에러도 알림도 없습니다. 고객은 "주문했는데 결제가 안 됐다"고 문의합니다. ScheduleToStart 를 5분으로 걸어 뒀다면 5분 만에
ActivityTaskTimedOut이 터져 워크플로우가 실패하고, 실패율 알림이 즉시 울렸을 것입니다. ScheduleToStart 는 성능 옵션이 아니라 관측 장치입니다.
10만 건을 처리하는 배치 액티비티가 있다고 합시다. StartToClose 는 2시간으로 잡았습니다. 30분쯤 지나 Worker 컨테이너가 OOM 으로 죽었습니다.
Temporal 서버는 이 사실을 모릅니다. 서버 입장에서 액티비티는 여전히 "시작됨" 상태이고, StartToClose 2시간이 지날 때까지 아무 일도 일어나지 않습니다. 1시간 30분을 낭비합니다.
Heartbeat 는 이 구멍을 메웁니다.
// [4-4] Heartbeat 를 보내는 액티비티
@Override
public String processBatch(String batchId, int totalCount) {
ActivityExecutionContext ctx = Activity.getExecutionContext();
for (int i = 0; i < totalCount; i++) {
processOne(batchId, i);
if (i % 1000 == 0) {
ctx.heartbeat(i); // 진행 상황을 details 로 함께 보낸다
log.info("배치 진행 {}/{}", i, totalCount);
}
}
return batchId + " DONE " + totalCount;
}ActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofHours(2))
.setHeartbeatTimeout(Duration.ofSeconds(30)) // 30초 안에 하트비트가 없으면 죽은 것으로 간주
.build();Worker 를 kill -9 로 죽입니다.
결과 (30초 뒤 히스토리)
{
"eventId": "7",
"eventTime": "2026-03-11T10:44:03.291Z",
"eventType": "EVENT_TYPE_ACTIVITY_TASK_TIMED_OUT",
"activityTaskTimedOutEventAttributes": {
"failure": {
"message": "activity Heartbeat timeout",
"timeoutFailureInfo": {
"timeoutType": "TIMEOUT_TYPE_HEARTBEAT",
"lastHeartbeatDetails": {
"payloads": [ { "metadata": { "encoding": "anNvbi9wbGFpbg==" }, "data": "NTAwMDA=" } ]
}
}
},
"scheduledEventId": "5",
"startedEventId": "6",
"retryState": "RETRY_STATE_IN_PROGRESS"
}
}1시간 30분 → 30초. 그리고 lastHeartbeatDetails 에 NTAwMDA= (base64 디코딩하면 50000)이 들어 있습니다.
재시도가 시작될 때, 액티비티는 마지막 하트비트 값을 꺼내 볼 수 있습니다.
// [4-4] 재시도 시 중단 지점부터 재개
@Override
public String processBatch(String batchId, int totalCount) {
ActivityExecutionContext ctx = Activity.getExecutionContext();
// 이전 시도가 남긴 진행 상황이 있으면 그 지점부터 시작한다
int start = ctx.getHeartbeatDetails(Integer.class).orElse(0);
log.info("배치 시작 batchId={} start={} total={}", batchId, start, totalCount);
for (int i = start; i < totalCount; i++) {
processOne(batchId, i);
if (i % 1000 == 0) {
ctx.heartbeat(i);
}
}
return batchId + " DONE " + totalCount;
}결과 (Worker 콘솔 — 첫 시도)
10:41:02.118 [Activity Executor taskQueue="ORDER_TASK_QUEUE": 1] INFO c.e.order.BatchActivityImpl - 배치 시작 batchId=B-77 start=0 total=100000
10:41:04.220 [Activity Executor taskQueue="ORDER_TASK_QUEUE": 1] INFO c.e.order.BatchActivityImpl - 배치 진행 0/100000
10:41:34.881 [Activity Executor taskQueue="ORDER_TASK_QUEUE": 1] INFO c.e.order.BatchActivityImpl - 배치 진행 25000/100000
10:42:31.402 [Activity Executor taskQueue="ORDER_TASK_QUEUE": 1] INFO c.e.order.BatchActivityImpl - 배치 진행 50000/100000여기서 kill -9.
결과 (Worker 재기동 후 — 두 번째 시도)
10:44:33.507 [Activity Executor taskQueue="ORDER_TASK_QUEUE": 1] INFO c.e.order.BatchActivityImpl - 배치 시작 batchId=B-77 start=50000 total=100000
10:44:34.610 [Activity Executor taskQueue="ORDER_TASK_QUEUE": 1] INFO c.e.order.BatchActivityImpl - 배치 진행 50000/100000
10:45:31.118 [Activity Executor taskQueue="ORDER_TASK_QUEUE": 1] INFO c.e.order.BatchActivityImpl - 배치 진행 75000/100000
10:46:29.774 [Activity Executor taskQueue="ORDER_TASK_QUEUE": 1] INFO c.e.order.BatchActivityImpl - 배치 완료 100000/1000005만 번째부터 재개했습니다. 처음부터 다시 했다면 약 3분 40초가 더 들었을 것을 1분 56초에 끝냈습니다.
💡 하트비트는 스로틀링됩니다
ctx.heartbeat()를 호출한다고 매번 서버로 gRPC 가 나가지는 않습니다. SDK 가 HeartbeatTimeout 의 80% 간격으로 묶어서 보냅니다. HeartbeatTimeout 이 30초면 실제 전송은 약 24초마다입니다. 그래서 1,000건마다 호출해도 서버 부하가 되지 않습니다. 반대로 HeartbeatTimeout 을 설정하지 않으면heartbeat()호출은 아무 효과가 없습니다. details 도 저장되지 않습니다. 이 조합을 자주 실수합니다.
워크플로우가 취소되거나(temporal workflow cancel) 액티비티 자체가 취소되면, 액티비티는 그것을 어떻게 알까요? 하트비트를 통해 알게 됩니다.
서버는 취소 요청을 받으면 그것을 기록해 두었다가, 액티비티가 다음 하트비트를 보내는 순간 응답에 "취소됨" 플래그를 실어 보냅니다. SDK 는 이를 받아 ctx.heartbeat() 호출 지점에서 ActivityCanceledException 을 던집니다.
@Override
public String processBatch(String batchId, int totalCount) {
ActivityExecutionContext ctx = Activity.getExecutionContext();
int start = ctx.getHeartbeatDetails(Integer.class).orElse(0);
for (int i = start; i < totalCount; i++) {
processOne(batchId, i);
if (i % 1000 == 0) {
try {
ctx.heartbeat(i);
} catch (ActivityCompletionException e) {
// ActivityCanceledException 은 ActivityCompletionException 의 하위 타입
log.info("취소 감지 — {}건 처리 후 정리하고 종료", i);
cleanupPartialWork(batchId, i);
throw e; // 반드시 다시 던져야 취소가 서버에 반영된다
}
}
}
return batchId + " DONE";
}temporal workflow cancel -w order-4001결과 (Worker 콘솔)
11:03:12.408 [Activity Executor taskQueue="ORDER_TASK_QUEUE": 1] INFO c.e.order.BatchActivityImpl - 배치 진행 33000/100000
11:03:36.117 [Activity Executor taskQueue="ORDER_TASK_QUEUE": 1] INFO c.e.order.BatchActivityImpl - 취소 감지 — 34000건 처리 후 정리하고 종료
11:03:36.120 [Activity Executor taskQueue="ORDER_TASK_QUEUE": 1] WARN i.t.i.a.ActivityTaskHandlerImpl - Activity failure. ActivityId=5, ActivityType=ProcessBatch, WorkflowId=order-4001
io.temporal.client.ActivityCanceledException: Activity cancelled by the workflow or the service
at io.temporal.internal.activity.ActivityExecutionContextImpl.doHeartBeat(ActivityExecutionContextImpl.java:141)
at io.temporal.internal.activity.ActivityExecutionContextImpl.heartbeat(ActivityExecutionContextImpl.java:103)
at com.example.order.BatchActivityImpl.processBatch(BatchActivityImpl.java:38)temporal workflow show -w order-4001결과
Progress:
ID Time Type
...
6 2026-03-11T11:02:41Z ActivityTaskStarted
7 2026-03-11T11:03:31Z WorkflowExecutionCancelRequested
8 2026-03-11T11:03:36Z ActivityTaskCanceled
9 2026-03-11T11:03:36Z WorkflowTaskScheduled
...
12 2026-03-11T11:03:36Z WorkflowExecutionCanceled
Result:
Status: CANCELED⚠️ 하트비트를 보내지 않는 액티비티는 취소를 감지할 수 없습니다 취소 전달 경로가 하트비트뿐이기 때문입니다. 하트비트 없는 액티비티는 취소 요청이 와도 끝까지 실행됩니다. 워크플로우는
Canceled로 닫히는데 액티비티는 뒤에서 계속 도는 상태가 됩니다. 오래 걸리는 액티비티에는 무조건 하트비트를 넣으세요. 취소 대응이 필요 없더라도 죽은 Worker 감지만으로 값어치를 합니다.
Temporal 이 보장하는 것은 at-least-once 입니다. exactly-once 가 아닙니다.
액티비티는 최소 한 번 실행됩니다. 두 번 실행될 수도 있습니다.
가장 무서운 시나리오는 "타임아웃으로 재시도했는데 사실 첫 시도가 성공했던" 경우입니다.
StartToClose = 10초, 결제 게이트웨이는 평소 2초
t=0s ActivityTaskScheduled
t=0s ActivityTaskStarted Worker A 가 집어감
t=0s Worker A → PG 로 HTTP POST /charges (39,000원)
t=2s PG 내부에서 결제 승인 완료. 카드사 승인 떨어짐. 💳 39,000원 결제됨
t=2s PG 가 200 응답을 보냄
↓
✘ 응답이 네트워크에서 유실됨 (LB 재시작 / TCP RST / Worker GC 정지)
↓
t=10s StartToClose 타임아웃. ActivityTaskTimedOut
t=11s 재시도. ActivityTaskScheduled (attempt=2)
t=11s Worker B 가 집어감
t=11s Worker B → PG 로 HTTP POST /charges (39,000원)
t=13s PG 가 이것을 **새 결제 요청**으로 처리. 💳 39,000원 또 결제됨
t=13s 200 응답 성공
t=13s ActivityTaskCompleted
워크플로우는 정상 완료. 히스토리에도 아무 이상 없음.
고객 카드에는 78,000원.히스토리를 아무리 봐도 이 사고는 보이지 않습니다. 이것이 "에러 없이 조용히 잘못 동작하는" 전형입니다.
// 워크플로우 코드
@Override
public String processOrder(OrderRequest req) {
// Run ID 를 시드로 하는 결정적 UUID. 리플레이해도 같은 값. (Step 03 의 3-7)
String idemKey = Workflow.randomUUID().toString();
String paymentId = payment.chargeIdempotent(req.orderId(), req.amount(), idemKey);
return req.orderId() + " COMPLETED";
}// 액티비티 구현
@Override
public String chargeIdempotent(String orderId, long amount, String idempotencyKey) {
HttpRequest req = HttpRequest.newBuilder()
.uri(URI.create("https://pg.example.com/v1/charges"))
.header("Idempotency-Key", idempotencyKey) // ← 여기
.POST(...)
.build();
// ...
}이제 타임라인이 이렇게 바뀝니다.
t=0s Worker A → PG (Idempotency-Key: aaa-111) → 💳 39,000원 결제
t=2s 응답 유실
t=10s 타임아웃
t=11s Worker B → PG (Idempotency-Key: aaa-111) ← 같은 키!
t=11s PG: "이 키는 이미 처리했다" → 저장된 첫 응답을 그대로 반환. 결제 안 함
t=11s ActivityTaskCompleted, paymentId = 첫 시도와 동일멱등키의 세 가지 조건을 정리합니다.
| 조건 | 이유 |
|---|---|
| 워크플로우에서 만든다 | 액티비티 안에서 만들면 재시도마다 새 키가 된다 |
Workflow.randomUUID() 를 쓴다 | UUID.randomUUID() 는 리플레이 때 값이 바뀐다 |
| 액티비티 인자로 전달한다 | 인자는 히스토리에 기록되어 모든 재시도가 같은 값을 받는다 |
액티비티 인자가 히스토리에 박혀 있으므로 재시도가 같은 키를 받는다는 것을 확인합니다.
temporal workflow show -w order-5001 --output json \
| jq '.events[] | select(.eventType=="EVENT_TYPE_ACTIVITY_TASK_SCHEDULED")
| .activityTaskScheduledEventAttributes.input.payloads[2].data' | base64 -d결과
"3f2a9c14-7b01-48d4-aef6-900c1d2f88a4"ActivityTaskScheduled 는 재시도할 때 새로 생기지 않습니다. 서버가 같은 이벤트의 입력으로 다시 디스패치할 뿐이라, 몇 번을 재시도해도 액티비티는 같은 세 번째 인자를 받습니다.
⚠️ 함정 — 멱등하지 않은 액티비티에 재시도를 켜 두면 결제가 두 번 됩니다 Temporal 의 기본 RetryOptions 는 무한 재시도입니다. 아무것도 설정하지 않으면 켜져 있습니다. 액티비티가 멱등하지 않다면 선택지는 셋입니다.
- 멱등하게 만든다 (권장). 외부 API 의 Idempotency-Key, DB 의
INSERT ... ON CONFLICT DO NOTHING, 유니크 제약.- 재시도를 끈다:
RetryOptions.newBuilder().setMaximumAttempts(1).build(). 대신 일시적 네트워크 오류에도 워크플로우가 실패합니다.- 실패를 non-retryable 로 분류한다:
setDoNotRetry("PaymentDuplicateError").실무에서는 1번 외에 답이 없습니다. "이 액티비티가 두 번 실행되면 무슨 일이 벌어지나" 를 모든 액티비티에 대해 답할 수 있어야 합니다.
💡 읽기 액티비티는 이미 멱등입니다 조회·계산·검증처럼 부수효과가 없는 액티비티는 신경 쓸 것이 없습니다. 멱등성이 문제가 되는 것은 쓰기뿐입니다: 결제, 이메일 발송, 재고 차감, 외부 시스템 등록.
일반 액티비티는 호출 한 번에 히스토리 이벤트가 최소 3개(Scheduled / Started / Completed) 생기고, Task Queue 를 거치는 서버 왕복이 발생합니다. "문자열 포맷팅", "간단한 검증", "로컬 캐시 조회" 같은 밀리초 단위 작업에는 과합니다.
Local Activity 는 Worker 프로세스 안에서 직접 실행되고, 히스토리에는 MarkerRecorded 하나만 남습니다.
private final ValidationActivity validation = Workflow.newLocalActivityStub(
ValidationActivity.class,
LocalActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofSeconds(5))
.build());10회 호출했을 때의 히스토리를 비교합니다.
일반 액티비티 10회
$ temporal workflow describe -w order-6001
History Length 36
History Size 8712Local Activity 10회
$ temporal workflow describe -w order-6002
History Length 16
History Size 3104temporal workflow show -w order-6002결과
Progress:
ID Time Type
1 2026-03-11T11:22:04Z WorkflowExecutionStarted
2 2026-03-11T11:22:04Z WorkflowTaskScheduled
3 2026-03-11T11:22:04Z WorkflowTaskStarted
4 2026-03-11T11:22:04Z WorkflowTaskCompleted
5 2026-03-11T11:22:04Z MarkerRecorded ← LocalActivity 1
6 2026-03-11T11:22:04Z MarkerRecorded ← LocalActivity 2
... (7~13 도 전부 MarkerRecorded)
14 2026-03-11T11:22:04Z MarkerRecorded ← LocalActivity 10
15 2026-03-11T11:22:04Z WorkflowTaskCompleted
16 2026-03-11T11:22:04Z WorkflowExecutionCompleted히스토리 36개 → 16개, 크기 8,712 → 3,104 바이트. 실행 시간은 1.84초 → 0.21초. 서버 왕복 10회가 0회로 줄었기 때문입니다.
| 항목 | Activity | Local Activity |
|---|---|---|
| 실행 위치 | Task Queue 를 거쳐 임의의 Worker | 워크플로우를 실행 중인 Worker 프로세스 안 |
| 히스토리 | Scheduled / Started / Completed (3개+) | MarkerRecorded 1개 |
| 서버 왕복 | 있음 (수십 ms) | 없음 |
| Task Queue 분리 | 가능 (GPU 전용 큐 등) | 불가 |
| 재시도 | 서버가 관리. 히스토리에 기록 | Worker 메모리에서. Workflow Task 타임아웃 안에서만 |
| Heartbeat | 지원 | 미지원 |
| 취소 | 지원 | 제한적 |
| 권장 실행 시간 | 제한 없음 | 수 초 이내 |
| Worker 재시작 시 | 재시도됨 | 처음부터 다시 (진행 상황 없음) |
| 워크플로우 타임아웃 영향 | 없음 | Workflow Task Timeout(기본 10초)에 갇힘 |
⚠️ 함정 — Local Activity 를 장기 작업에 쓰면 Local Activity 는 Workflow Task 안에서 실행됩니다. Workflow Task Timeout 기본값은 10초입니다. Local Activity 가 그보다 오래 걸리면 SDK 가 내부적으로 Workflow Task 를 heartbeat 하며 연장을 시도하지만, 한계가 있습니다. 결과
io.temporal.internal.statemachines.LocalActivityStateMachine: Local activity exceeded workflow task timeout. Rescheduling as a new workflow task.결국 히스토리가 오히려 지저분해지고, Worker 가 죽으면 진행 상황이 전부 사라집니다. 판단 기준: 1초 이내에 끝나고, 재시도해도 부담 없고, 외부 시스템을 안 건드리면 Local Activity. 그 외에는 전부 일반 Activity.
워크플로우는 클래스를 등록하지만 액티비티는 인스턴스를 등록합니다.
Worker worker = factory.newWorker(TASK_QUEUE);
// 워크플로우: 클래스 — SDK 가 실행마다 새 인스턴스를 만든다
worker.registerWorkflowImplementationTypes(OrderWorkflowImpl.class);
// 액티비티: 인스턴스 — Worker 가 이 하나를 계속 재사용한다
worker.registerActivitiesImplementations(
new PaymentActivityImpl(paymentClient),
new InventoryActivityImpl(inventoryRepo),
new ShippingActivityImpl(shippingClient),
new NotificationActivityImpl(mailer));
factory.start();여기서 결론이 하나 나옵니다.
액티비티 구현체는 스레드 안전해야 합니다.
한 Worker 는 기본적으로 액티비티를 동시에 200개까지 실행합니다(maxConcurrentActivityExecutionSize). 그 200개 스레드가 같은 인스턴스의 메서드를 호출합니다.
// ✘ 절대 하지 말 것
public class PaymentActivityImpl implements PaymentActivity {
private String currentOrderId; // 인스턴스 필드에 요청별 상태!
private long total;
@Override
public String charge(String orderId, long amount) {
this.currentOrderId = orderId; // 다른 스레드가 덮어쓴다
this.total += amount; // 경쟁 조건
// ... 200 스레드가 이 필드를 서로 짓밟는다
return doCharge(this.currentOrderId, amount); // 남의 주문을 결제할 수 있다
}
}// ✔ 무상태로. 공유 필드는 불변이거나 스레드 안전한 것만.
public class PaymentActivityImpl implements PaymentActivity {
private final HttpClient http; // HttpClient 는 스레드 안전
private final PaymentRepository repo; // 스프링 리포지터리도 스레드 안전
public PaymentActivityImpl(HttpClient http, PaymentRepository repo) {
this.http = http;
this.repo = repo;
}
@Override
public String charge(String orderId, long amount) {
// 모든 상태는 지역 변수와 파라미터로만
String requestId = UUID.randomUUID().toString();
return doCharge(http, orderId, amount, requestId);
}
}동시 실행 수와 처리율은 WorkerOptions 로 조절합니다.
WorkerOptions options = WorkerOptions.newBuilder()
.setMaxConcurrentActivityExecutionSize(50) // 기본 200
.setMaxConcurrentWorkflowTaskExecutionSize(100) // 기본 200
.setMaxWorkerActivitiesPerSecond(20.0) // Worker 당 초당 액티비티 수 제한
.setMaxTaskQueueActivitiesPerSecond(100.0) // Task Queue 전체 초당 제한
.build();
Worker worker = factory.newWorker(TASK_QUEUE, options);결과 (temporal task-queue describe --task-queue ORDER_TASK_QUEUE --task-queue-type activity)
BuildId TaskQueueType Pollers BacklogCount ApproximateBacklogAge
ACTIVITY 4 0 0s
Pollers:
Identity LastAccessTime RatePerSecond
41822@worker-7d9f4c6b8-xk2mp 2026-03-11T11:31:02Z 100000
41822@worker-7d9f4c6b8-p4nzq 2026-03-11T11:31:02Z 100000💡 실무 팁 — 무거운 액티비티는 Task Queue 를 분리하세요 이미지 변환이나 ML 추론처럼 CPU·메모리를 많이 먹는 액티비티가 결제 액티비티와 같은 Worker 에서 돌면, 배치 하나가 결제 전체를 굶깁니다.
// 워크플로우 코드에서 액티비티별로 다른 Task Queue 지정 private final MediaActivity media = Workflow.newActivityStub( MediaActivity.class, ActivityOptions.newBuilder() .setTaskQueue("MEDIA_TASK_QUEUE") // ← 전용 큐 .setStartToCloseTimeout(Duration.ofMinutes(30)) .setHeartbeatTimeout(Duration.ofSeconds(30)) .build());그리고 GPU 인스턴스에
MEDIA_TASK_QUEUE만 폴링하는 Worker 를 따로 띄웁니다. 워크플로우 코드는 한 줄도 안 바뀝니다.
| 개념 | 핵심 |
|---|---|
@ActivityInterface | 메서드 여러 개 가능. 구현체는 평범한 POJO — 제약 없음 |
@ActivityMethod(name=...) | 액티비티 타입명은 히스토리에 문자열로 저장. 메서드명 변경 시 못 박을 것 |
ActivityOptions | 워크플로우 코드 안에서 만든다. 히스토리에 기록됨 |
| ScheduleToStart | 큐 대기 시간. 기본 무한. 재시도 안 됨. Worker 부족·잘못된 큐를 잡는 관측 장치 |
| StartToClose | 한 번의 시도 시간. 재시도됨. 가장 중요한 타임아웃 |
| ScheduleToClose | 재시도 포함 총 시간. 비즈니스 데드라인 |
| Heartbeat | 하트비트 간격. 조용히 죽은 Worker 를 초 단위로 감지 |
| 필수 규칙 | ScheduleToClose 또는 StartToClose 중 하나는 반드시. 없으면 워크플로우 FAILED |
| 타임아웃 이벤트 | ActivityTaskTimedOut 의 timeoutType 으로 어느 타임아웃인지 구분 |
heartbeat(details) | 진행 상황 저장. 재시도 시 getHeartbeatDetails() 로 이어서 실행 |
| 하트비트 스로틀링 | HeartbeatTimeout 의 80% 간격으로 묶어 전송. HeartbeatTimeout 미설정 시 무효 |
| 취소 감지 | heartbeat() 가 ActivityCanceledException 을 던진다. 하트비트 없으면 취소를 모른다 |
| at-least-once | 액티비티는 두 번 실행될 수 있다. 응답 유실 후 재시도가 전형적 |
| 멱등키 | 워크플로우에서 Workflow.randomUUID() 로 만들어 액티비티 인자로 전달 |
| Local Activity | MarkerRecorded 1개만. 히스토리 36→16, 1.84초→0.21초. 단 하트비트·장기실행 불가 |
| 액티비티 등록 | 인스턴스를 등록 → Worker 당 하나를 200 스레드가 공유 → 스레드 안전 필수 |
| Task Queue 분리 | 무거운 액티비티는 전용 큐로. setTaskQueue() 한 줄 |
Exercise.java 에 6문제가 있습니다. 정답은 Solution.java.
액티비티가 타임아웃으로 실패하면 Temporal 이 재시도한다는 것을 봤습니다. 그런데 그 재시도는 몇 번까지, 얼마 간격으로 일어날까요. 재시도해도 소용없는 실패(잘못된 카드번호)와 재시도해야 하는 실패(네트워크 오류)는 어떻게 구분할까요. 다음 스텝에서 RetryOptions, ApplicationFailure, non-retryable 오류 분류를 다룹니다.
이 스텝은 Java 파일 세 개로 진행합니다. Practice.java 로 4-1 ~ 4-8 의 예제를 재현하고, Exercise.java 의 6문제를 푼 뒤, Solution.java 로 대조합니다. Practice 의 타임아웃 실습은 의도적으로 실패하는 워크플로우를 여러 개 돌리므로, 실행 후 temporal workflow list 에 FAILED 가 쌓이는 것이 정상입니다.
본문의 모든 예제를 절 번호 주석과 함께 담은 실습 파일입니다.
[4-3] 은 --timeout 인자로 실행합니다. SlowPaymentActivityImpl 이 10초를 자는 동안 StartToClose 3초가 터지도록 되어 있어, 실행 후 temporal workflow show -w order-3002 --output json 으로 TIMEOUT_TYPE_START_TO_CLOSE 를 확인할 수 있습니다.[4-3] 의 ScheduleToStart 실습은 Worker 를 띄우지 않고 워크플로우만 시작해야 재현됩니다. --no-worker 인자를 주면 클라이언트만 떠서 액티비티가 큐에 쌓이고, 10초 뒤 TIMEOUT_TYPE_SCHEDULE_TO_START 가 터집니다.[4-4] 의 BatchActivityImpl 이 10만 건 재개 실습의 핵심입니다. KILL_AT 상수(기본 50000)에 도달하면 스스로 System.exit(137) 로 죽어 kill -9 를 흉내 냅니다. Worker 를 다시 띄우면 start=50000 부터 재개되는 로그를 볼 수 있습니다.[4-6] 의 NonIdempotentPaymentActivityImpl 은 호출될 때마다 잔액을 실제로 깎는 가짜 원장을 들고 있습니다. --double-charge 로 실행하면 응답 유실을 흉내 내 78,000원이 빠지는 것을 콘솔에서 볼 수 있습니다. 같은 파일의 IdempotentPaymentActivityImpl 이 이를 막습니다.[4-7] 은 LocalActivityWorkflowImpl 과 RegularActivityWorkflowImpl 을 둘 다 담고 있습니다. 각각 실행한 뒤 temporal workflow describe 의 History Length 와 History Size 를 비교하세요. 36 vs 16, 8712 vs 3104 가 나옵니다.[4-8] 의 UnsafePaymentActivityImpl 은 인스턴스 필드에 요청별 상태를 담은 반면교사입니다. --race 로 실행하면 20개 워크플로우를 동시에 던져 다른 주문의 orderId 로 결제되는 로그를 재현합니다.package com.example.order;
// =====================================================================================
// Step 04 — 액티비티 : Practice
//
// 실행 방법
// ./gradlew run -PmainClass=com.example.order.Practice
// ./gradlew run -PmainClass=com.example.order.Practice --args="--timeout" [4-3]
// ./gradlew run -PmainClass=com.example.order.Practice --args="--no-worker" [4-3] S2S
// ./gradlew run -PmainClass=com.example.order.Practice --args="--heartbeat" [4-4]
// ./gradlew run -PmainClass=com.example.order.Practice --args="--double-charge" [4-6]
// ./gradlew run -PmainClass=com.example.order.Practice --args="--local" [4-7]
// ./gradlew run -PmainClass=com.example.order.Practice --args="--race" [4-8]
//
// 사전 조건
// docker compose up -d (Temporal Server 1.22.4, gRPC 127.0.0.1:7233)
//
// 주의
// 타임아웃 실습은 **의도적으로 실패하는** 워크플로우를 돌립니다.
// 실행 후 temporal workflow list 에 FAILED 가 쌓이는 것이 정상입니다.
// =====================================================================================
import io.temporal.activity.Activity;
import io.temporal.activity.ActivityExecutionContext;
import io.temporal.activity.ActivityInterface;
import io.temporal.activity.ActivityMethod;
import io.temporal.activity.ActivityOptions;
import io.temporal.activity.LocalActivityOptions;
import io.temporal.client.ActivityCompletionException;
import io.temporal.client.WorkflowClient;
import io.temporal.client.WorkflowOptions;
import io.temporal.common.RetryOptions;
import io.temporal.failure.ApplicationFailure;
import io.temporal.serviceclient.WorkflowServiceStubs;
import io.temporal.worker.Worker;
import io.temporal.worker.WorkerFactory;
import io.temporal.worker.WorkerOptions;
import io.temporal.workflow.Workflow;
import io.temporal.workflow.WorkflowInterface;
import io.temporal.workflow.WorkflowMethod;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.time.Duration;
import java.time.Instant;
import java.time.ZoneId;
import java.time.format.DateTimeFormatter;
import java.util.Map;
import java.util.UUID;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.atomic.AtomicLong;
public class Practice {
public static final String TASK_QUEUE = "ORDER_TASK_QUEUE";
public static final String MEDIA_TASK_QUEUE = "MEDIA_TASK_QUEUE";
public record OrderRequest(
String orderId, String customerId, String sku, int qty, long amount, String address) {}
@WorkflowInterface
public interface OrderWorkflow {
@WorkflowMethod String processOrder(OrderRequest req);
}
// =================================================================================
// [4-1] 액티비티 인터페이스 — 메서드 여러 개 가능. @ActivityMethod 는 생략 가능하다.
//
// ★ @ActivityMethod(name = "...") 으로 타입명을 못 박는 이유:
// 액티비티 타입 이름은 히스토리에 문자열로 저장된다.
// "activityType": { "name": "Charge" }
// 자바 메서드명을 chargeWithFee 로 리팩터링하면, 히스토리에 "Charge" 가 적힌
// 실행 중인 워크플로우를 리플레이할 때 SDK 가 타입을 못 찾는다.
// ApplicationFailure: Activity Type "Charge" is not registered with a worker.
// 이름을 못 박아 두면 메서드명은 자유롭게 바꿀 수 있다. (Step 10 버저닝)
// =================================================================================
@ActivityInterface
public interface PaymentActivity {
@ActivityMethod(name = "Charge")
String charge(String orderId, long amount);
@ActivityMethod(name = "ChargeIdempotent")
String chargeIdempotent(String orderId, long amount, String idempotencyKey);
@ActivityMethod(name = "Refund")
void refund(String paymentId);
}
@ActivityInterface
public interface BatchActivity {
@ActivityMethod(name = "ProcessBatch")
String processBatch(String batchId, int totalCount);
}
@ActivityInterface
public interface ValidationActivity {
@ActivityMethod(name = "NormalizeAddress")
String normalizeAddress(String address);
}
// =================================================================================
// [4-1] 액티비티 구현체 — 평범한 POJO. 벽시계·랜덤·네트워크·스레드 전부 허용된다.
// =================================================================================
public static class PaymentActivityImpl implements PaymentActivity {
private static final Logger log = LoggerFactory.getLogger(PaymentActivityImpl.class);
@Override
public String charge(String orderId, long amount) {
// 벽시계 — 액티비티에서는 문제없다
long start = System.currentTimeMillis();
// 랜덤 — 문제없다
String requestId = UUID.randomUUID().toString();
log.info("charge orderId={} amount={} requestId={}", orderId, amount, requestId);
sleepQuietly(300); // 외부 PG 호출을 흉내
log.info("PG 응답 {}ms", System.currentTimeMillis() - start);
return "PAY-" + requestId.substring(0, 8);
}
@Override
public String chargeIdempotent(String orderId, long amount, String idempotencyKey) {
log.info("chargeIdempotent orderId={} amount={} key={}", orderId, amount, idempotencyKey);
sleepQuietly(300);
return "PAY-" + idempotencyKey.substring(0, 8);
}
@Override
public void refund(String paymentId) {
log.info("refund paymentId={}", paymentId);
}
}
// ---------------------------------------------------------------------------------
// [4-3] 일부러 느린 결제 액티비티. StartToClose 3초를 터뜨리는 데 쓴다.
// ---------------------------------------------------------------------------------
public static class SlowPaymentActivityImpl implements PaymentActivity {
private static final Logger log = LoggerFactory.getLogger(SlowPaymentActivityImpl.class);
@Override
public String charge(String orderId, long amount) {
log.info("느린 결제 시작 orderId={} (10초 걸린다)", orderId);
sleepQuietly(10_000); // StartToClose 3초 < 10초 → 타임아웃
log.info("느린 결제 완료 orderId={}", orderId);
return "PAY-SLOW";
}
@Override
public String chargeIdempotent(String orderId, long amount, String key) {
return charge(orderId, amount);
}
@Override
public void refund(String paymentId) {}
}
// =================================================================================
// [4-3] 타임아웃 4종. 각 워크플로우가 하나씩을 터뜨린다.
// =================================================================================
// (a) 타임아웃을 하나도 안 준다 → 워크플로우가 즉시 FAILED
// IllegalStateException: Either ScheduleToCloseTimeout or StartToCloseTimeout is required
public static class NoTimeoutWorkflowImpl implements OrderWorkflow {
private final PaymentActivity payment = Workflow.newActivityStub(
PaymentActivity.class,
ActivityOptions.newBuilder().build()); // ✘
@Override
public String processOrder(OrderRequest req) {
payment.charge(req.orderId(), req.amount());
return req.orderId() + " COMPLETED";
}
}
// (b) StartToClose 3초 — 액티비티는 10초 걸린다 → TIMEOUT_TYPE_START_TO_CLOSE
public static class StartToCloseWorkflowImpl implements OrderWorkflow {
private final PaymentActivity payment = Workflow.newActivityStub(
PaymentActivity.class,
ActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofSeconds(3))
.setRetryOptions(RetryOptions.newBuilder().setMaximumAttempts(2).build())
.build());
@Override
public String processOrder(OrderRequest req) {
payment.charge(req.orderId(), req.amount());
return req.orderId() + " COMPLETED";
}
}
// (c) ScheduleToStart 10초 — Worker 를 안 띄우면 큐에서 대기하다 터진다
// → TIMEOUT_TYPE_SCHEDULE_TO_START, startedEventId=0, RETRY_STATE_NON_RETRYABLE_FAILURE
public static class ScheduleToStartWorkflowImpl implements OrderWorkflow {
private final PaymentActivity payment = Workflow.newActivityStub(
PaymentActivity.class,
ActivityOptions.newBuilder()
.setScheduleToStartTimeout(Duration.ofSeconds(10))
.setStartToCloseTimeout(Duration.ofSeconds(30))
.build());
@Override
public String processOrder(OrderRequest req) {
payment.charge(req.orderId(), req.amount());
return req.orderId() + " COMPLETED";
}
}
// (d) ScheduleToClose 20초 — 재시도를 계속하다 전체 예산이 소진된다
// → TIMEOUT_TYPE_SCHEDULE_TO_CLOSE, failure.cause 에 직전 실패가 중첩된다
public static class ScheduleToCloseWorkflowImpl implements OrderWorkflow {
private final PaymentActivity payment = Workflow.newActivityStub(
PaymentActivity.class,
ActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofSeconds(3))
.setScheduleToCloseTimeout(Duration.ofSeconds(20))
.build());
@Override
public String processOrder(OrderRequest req) {
payment.charge(req.orderId(), req.amount());
return req.orderId() + " COMPLETED";
}
}
// =================================================================================
// [4-4] Heartbeat — 10만 건 배치를 5만 건에서 죽였다가 이어서 한다.
//
// KILL_AT 에 도달하면 System.exit(137) 로 스스로 죽는다 (kill -9 흉내).
// Worker 를 다시 띄우면 getHeartbeatDetails() 가 50000 을 돌려주어 거기서부터 재개한다.
//
// ★ HeartbeatTimeout 을 설정하지 않으면 heartbeat() 호출은 아무 효과가 없다.
// details 도 저장되지 않는다. 이 조합을 자주 실수한다.
// =================================================================================
public static class BatchActivityImpl implements BatchActivity {
private static final Logger log = LoggerFactory.getLogger(BatchActivityImpl.class);
private static final int KILL_AT = 50_000;
private static volatile boolean killEnabled = false;
@Override
public String processBatch(String batchId, int totalCount) {
ActivityExecutionContext ctx = Activity.getExecutionContext();
// ★ 반복문 밖에서 딱 한 번 호출한다. 이 값은 액티비티 시도가 시작될 때
// 서버가 함께 내려 준 것이라 반복 호출할 이유가 없다.
int start = ctx.getHeartbeatDetails(Integer.class).orElse(0);
int attempt = ctx.getInfo().getAttempt();
log.info("배치 시작 batchId={} start={} total={} attempt={}",
batchId, start, totalCount, attempt);
for (int i = start; i < totalCount; i++) {
processOne(batchId, i);
if (i % 1000 == 0) {
try {
// 하트비트는 SDK 가 HeartbeatTimeout 의 80% 간격으로 묶어 보낸다.
// 1,000건마다 호출해도 서버 부하가 되지 않는다.
ctx.heartbeat(i);
} catch (ActivityCompletionException e) {
// [4-5] 취소 감지 — ActivityCanceledException 이 여기로 온다
log.info("취소 감지 — {}건 처리 후 정리하고 종료", i);
cleanupPartialWork(batchId, i);
throw e; // 반드시 다시 던져야 취소가 서버에 반영된다
}
if (i % 25_000 == 0) {
log.info("배치 진행 {}/{}", i, totalCount);
}
if (killEnabled && i >= KILL_AT) {
log.error("=== {}건 지점에서 Worker 강제 종료 (kill -9 흉내) ===", i);
log.error("=== Worker 를 다시 띄우면 {} 부터 재개됩니다 ===", i);
Runtime.getRuntime().halt(137);
}
}
}
log.info("배치 완료 {}/{}", totalCount, totalCount);
return batchId + " DONE " + totalCount;
}
private void processOne(String batchId, int i) {
// 실제 작업 흉내. 10만 건에 약 4분.
try {
Thread.sleep(0, 200_000);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
private void cleanupPartialWork(String batchId, int processed) {
log.info("부분 작업 정리 batchId={} processed={}", batchId, processed);
}
static void enableKill() {
killEnabled = true;
}
}
public static class HeartbeatWorkflowImpl implements OrderWorkflow {
private final BatchActivity batch = Workflow.newActivityStub(
BatchActivity.class,
ActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofHours(2))
// 30초 안에 하트비트가 없으면 Worker 가 죽은 것으로 간주하고 재시도한다.
// 이게 없으면 StartToClose 2시간이 다 지나야 알아챈다.
.setHeartbeatTimeout(Duration.ofSeconds(30))
.build());
@Override
public String processOrder(OrderRequest req) {
return batch.processBatch("B-" + req.orderId(), 100_000);
}
}
// =================================================================================
// [4-6] 멱등성 — at-least-once 라는 계약
//
// 아래 두 구현을 --double-charge 로 나란히 돌려 보면 차이가 명확하다.
// NonIdempotent 는 응답 유실 후 재시도에서 39,000원을 또 깎아 78,000원이 된다.
// =================================================================================
public static class NonIdempotentPaymentActivityImpl implements PaymentActivity {
private static final Logger log = LoggerFactory.getLogger(NonIdempotentPaymentActivityImpl.class);
// 가짜 원장. 실제로는 PG 서버의 상태다.
static final AtomicLong LEDGER = new AtomicLong(0);
private static final ConcurrentHashMap<String, Integer> ATTEMPTS = new ConcurrentHashMap<>();
@Override
public String charge(String orderId, long amount) {
int n = ATTEMPTS.merge(orderId, 1, Integer::sum);
// ★ 실제로 돈이 빠진다
long total = LEDGER.addAndGet(amount);
log.warn("💳 결제 실행 orderId={} amount={} (attempt={}) 누적={}", orderId, amount, n, total);
if (n == 1) {
// 첫 시도: 결제는 성공했지만 응답이 유실된다 (LB 재시작 / TCP RST / GC 정지)
log.error("✘ 응답 유실 시뮬레이션 — 결제는 됐지만 워커가 결과를 못 받는다");
sleepQuietly(15_000); // StartToClose 를 넘겨 타임아웃 유발
}
return "PAY-" + orderId + "-" + n;
}
@Override
public String chargeIdempotent(String orderId, long amount, String key) {
return charge(orderId, amount);
}
@Override
public void refund(String paymentId) {}
}
public static class IdempotentPaymentActivityImpl implements PaymentActivity {
private static final Logger log = LoggerFactory.getLogger(IdempotentPaymentActivityImpl.class);
static final AtomicLong LEDGER = new AtomicLong(0);
// PG 서버가 들고 있는 멱등키 → 결과 맵. 실제로는 PG 쪽 DB 다.
private static final ConcurrentHashMap<String, String> PROCESSED = new ConcurrentHashMap<>();
private static final ConcurrentHashMap<String, Integer> ATTEMPTS = new ConcurrentHashMap<>();
@Override
public String chargeIdempotent(String orderId, long amount, String idempotencyKey) {
int n = ATTEMPTS.merge(orderId, 1, Integer::sum);
// ★ 이미 처리한 키면 저장된 결과를 그대로 돌려준다. 돈이 빠지지 않는다.
String existing = PROCESSED.get(idempotencyKey);
if (existing != null) {
log.info("✔ 멱등키 중복 감지 key={} — 저장된 결과 반환 (결제 안 함) 누적={}",
idempotencyKey, LEDGER.get());
return existing;
}
long total = LEDGER.addAndGet(amount);
String paymentId = "PAY-" + idempotencyKey.substring(0, 8);
PROCESSED.put(idempotencyKey, paymentId);
log.warn("💳 결제 실행 orderId={} amount={} key={} (attempt={}) 누적={}",
orderId, amount, idempotencyKey, n, total);
if (n == 1) {
log.error("✘ 응답 유실 시뮬레이션");
sleepQuietly(15_000);
}
return paymentId;
}
@Override
public String charge(String orderId, long amount) {
throw ApplicationFailure.newNonRetryableFailure(
"멱등키 없는 charge 는 쓰지 마세요", "UnsafeChargeError");
}
@Override
public void refund(String paymentId) {}
}
// 멱등키를 워크플로우에서 만들어 액티비티 인자로 넘긴다.
// ★ Workflow.randomUUID() 여야 한다. UUID.randomUUID() 는 리플레이 때 값이 바뀐다.
public static class IdempotentWorkflowImpl implements OrderWorkflow {
private static final Logger log = Workflow.getLogger(IdempotentWorkflowImpl.class);
private final PaymentActivity payment = Workflow.newActivityStub(
PaymentActivity.class,
ActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofSeconds(5))
.setRetryOptions(RetryOptions.newBuilder().setMaximumAttempts(3).build())
.build());
@Override
public String processOrder(OrderRequest req) {
String idemKey = Workflow.randomUUID().toString();
log.info("[{}] 멱등키 발급 {}", req.orderId(), idemKey);
String paymentId = payment.chargeIdempotent(req.orderId(), req.amount(), idemKey);
log.info("[{}] 결제 완료 {}", req.orderId(), paymentId);
return req.orderId() + " COMPLETED";
}
}
public static class NonIdempotentWorkflowImpl implements OrderWorkflow {
private final PaymentActivity payment = Workflow.newActivityStub(
PaymentActivity.class,
ActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofSeconds(5))
.setRetryOptions(RetryOptions.newBuilder().setMaximumAttempts(3).build())
.build());
@Override
public String processOrder(OrderRequest req) {
payment.charge(req.orderId(), req.amount());
return req.orderId() + " COMPLETED";
}
}
// =================================================================================
// [4-7] Local Activity vs Activity — 히스토리 비용 비교
//
// 같은 액티비티를 10회 호출한다.
// Regular : History Length 36, History Size 8712, 1.84초
// Local : History Length 16, History Size 3104, 0.21초
// =================================================================================
public static class ValidationActivityImpl implements ValidationActivity {
private static final Logger log = LoggerFactory.getLogger(ValidationActivityImpl.class);
@Override
public String normalizeAddress(String address) {
// 순수 문자열 처리. 밀리초 단위. 외부 시스템을 안 건드린다 → Local Activity 후보
return address.trim().replaceAll("\\s+", " ").toUpperCase();
}
}
public static class RegularActivityWorkflowImpl implements OrderWorkflow {
private final ValidationActivity validation = Workflow.newActivityStub(
ValidationActivity.class,
ActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofSeconds(5))
.build());
@Override
public String processOrder(OrderRequest req) {
String addr = req.address();
for (int i = 0; i < 10; i++) {
addr = validation.normalizeAddress(addr); // 매번 서버 왕복 + 이벤트 3개
}
return req.orderId() + " REGULAR " + addr;
}
}
public static class LocalActivityWorkflowImpl implements OrderWorkflow {
private final ValidationActivity validation = Workflow.newLocalActivityStub(
ValidationActivity.class,
LocalActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofSeconds(5))
.build());
@Override
public String processOrder(OrderRequest req) {
String addr = req.address();
for (int i = 0; i < 10; i++) {
addr = validation.normalizeAddress(addr); // 서버 왕복 없음, MarkerRecorded 1개
}
return req.orderId() + " LOCAL " + addr;
}
}
// =================================================================================
// [4-8] 액티비티는 **인스턴스**로 등록된다 → Worker 당 하나를 200 스레드가 공유한다.
// 따라서 스레드 안전해야 한다.
// =================================================================================
public static class UnsafePaymentActivityImpl implements PaymentActivity {
private static final Logger log = LoggerFactory.getLogger(UnsafePaymentActivityImpl.class);
// ✘ 인스턴스 필드에 요청별 상태! 200 스레드가 서로 짓밟는다.
private String currentOrderId;
private long total;
@Override
public String charge(String orderId, long amount) {
this.currentOrderId = orderId;
sleepQuietly(50); // 다른 스레드가 끼어들 틈
this.total += amount; // 경쟁 조건 (원자적이지 않음)
if (!orderId.equals(this.currentOrderId)) {
log.error("★★★ 오염 감지! 요청={} 인데 필드={} — 남의 주문을 결제한다",
orderId, this.currentOrderId);
}
return "PAY-" + this.currentOrderId; // 남의 orderId 가 나갈 수 있다
}
@Override
public String chargeIdempotent(String orderId, long amount, String key) {
return charge(orderId, amount);
}
@Override
public void refund(String paymentId) {}
}
// ✔ 무상태로. 공유 필드는 불변이거나 스레드 안전한 것만.
public static class SafePaymentActivityImpl implements PaymentActivity {
private static final Logger log = LoggerFactory.getLogger(SafePaymentActivityImpl.class);
// DateTimeFormatter 는 불변이고 스레드 안전하다. SimpleDateFormat 은 아니다.
private final DateTimeFormatter fmt =
DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss").withZone(ZoneId.of("Asia/Seoul"));
@Override
public String charge(String orderId, long amount) {
// 모든 상태는 지역 변수와 파라미터로만
String requestId = UUID.randomUUID().toString();
String at = fmt.format(Instant.now());
log.info("charge orderId={} amount={} at={} requestId={}", orderId, amount, at, requestId);
sleepQuietly(50);
return "PAY-" + orderId + "-" + requestId.substring(0, 8);
}
@Override
public String chargeIdempotent(String orderId, long amount, String key) {
return "PAY-" + key.substring(0, 8);
}
@Override
public void refund(String paymentId) {}
}
private static void sleepQuietly(long ms) {
try {
Thread.sleep(ms);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
}
}
// =================================================================================
// 실행 진입점
// =================================================================================
public static void main(String[] args) throws Exception {
String mode = args.length > 0 ? args[0] : "--default";
WorkflowServiceStubs service = WorkflowServiceStubs.newLocalServiceStubs();
WorkflowClient client = WorkflowClient.newInstance(service);
WorkerFactory factory = WorkerFactory.newInstance(client);
// [4-8] Worker 튜닝
WorkerOptions workerOptions = WorkerOptions.newBuilder()
.setMaxConcurrentActivityExecutionSize(50) // 기본 200
.setMaxConcurrentWorkflowTaskExecutionSize(100) // 기본 200
.setMaxConcurrentLocalActivityExecutionSize(100) // 기본 200
.build();
Worker worker = factory.newWorker(TASK_QUEUE, workerOptions);
switch (mode) {
case "--timeout" -> {
// [4-3] StartToClose 3초 vs 10초짜리 액티비티
worker.registerWorkflowImplementationTypes(StartToCloseWorkflowImpl.class);
worker.registerActivitiesImplementations(new SlowPaymentActivityImpl());
factory.start();
start(client, "3002", "StartToClose 3초 — 10초짜리 액티비티가 터진다");
System.out.println("확인: temporal workflow show -w order-3002 --output json \\");
System.out.println(" | jq '.events[] | select(.eventType==\"EVENT_TYPE_ACTIVITY_TASK_TIMED_OUT\")'");
}
case "--no-worker" -> {
// [4-3] ScheduleToStart — Worker 를 띄우지 않는다
System.out.println("Worker 를 띄우지 않습니다. 액티비티가 큐에서 대기하다 터집니다.");
start(client, "3003", "ScheduleToStart 10초 — Worker 없음");
System.out.println("확인: temporal task-queue describe --task-queue " + TASK_QUEUE
+ " --task-queue-type activity");
Thread.sleep(15_000);
System.out.println("확인: temporal workflow show -w order-3003");
System.exit(0);
}
case "--heartbeat" -> {
// [4-4] 10만 건 배치 재개
BatchActivityImpl.enableKill();
worker.registerWorkflowImplementationTypes(HeartbeatWorkflowImpl.class);
worker.registerActivitiesImplementations(new BatchActivityImpl());
factory.start();
start(client, "4001", "10만 건 배치 — 5만 건에서 스스로 죽는다");
System.out.println();
System.out.println("=== [4-4] 재현 절차 ===");
System.out.println("1) 5만 건 지점에서 프로세스가 halt(137) 로 죽습니다.");
System.out.println("2) 30초 뒤 HeartbeatTimeout 이 터집니다:");
System.out.println(" temporal workflow show -w order-4001 --output json | \\");
System.out.println(" jq '.events[] | select(.eventType==\"EVENT_TYPE_ACTIVITY_TASK_TIMED_OUT\")'");
System.out.println(" → timeoutType: TIMEOUT_TYPE_HEARTBEAT, lastHeartbeatDetails: 50000");
System.out.println("3) 이 명령을 --args=\"--heartbeat-resume\" 으로 다시 실행하면");
System.out.println(" start=50000 부터 재개하는 로그가 보입니다.");
}
case "--heartbeat-resume" -> {
// 죽지 않는 Worker 로 재기동 → 5만 건부터 재개
worker.registerWorkflowImplementationTypes(HeartbeatWorkflowImpl.class);
worker.registerActivitiesImplementations(new BatchActivityImpl());
factory.start();
System.out.println("Worker 재기동. order-4001 이 5만 건부터 재개됩니다.");
}
case "--double-charge" -> {
// [4-6] 멱등하지 않은 결제 → 78,000원
worker.registerWorkflowImplementationTypes(NonIdempotentWorkflowImpl.class);
worker.registerActivitiesImplementations(new NonIdempotentPaymentActivityImpl());
factory.start();
start(client, "5001", "멱등하지 않은 결제 — 39,000원이 두 번 빠진다");
Thread.sleep(30_000);
System.out.println();
System.out.println("★ 원장 잔액: " + NonIdempotentPaymentActivityImpl.LEDGER.get() + "원");
System.out.println(" 기대값 39000, 실제값 78000 → 결제가 두 번 됐습니다.");
System.out.println(" 히스토리에는 아무 이상이 없습니다. temporal workflow show -w order-5001");
System.exit(0);
}
case "--idempotent" -> {
// [4-6] 멱등키로 막는다
worker.registerWorkflowImplementationTypes(IdempotentWorkflowImpl.class);
worker.registerActivitiesImplementations(new IdempotentPaymentActivityImpl());
factory.start();
start(client, "5002", "멱등키 결제 — 재시도해도 39,000원");
Thread.sleep(30_000);
System.out.println();
System.out.println("★ 원장 잔액: " + IdempotentPaymentActivityImpl.LEDGER.get() + "원");
System.exit(0);
}
case "--local" -> {
// [4-7] Local Activity vs Activity
worker.registerWorkflowImplementationTypes(
RegularActivityWorkflowImpl.class, LocalActivityWorkflowImpl.class);
worker.registerActivitiesImplementations(new ValidationActivityImpl());
factory.start();
startTyped(client, RegularActivityWorkflowImpl.class, "6001");
startTyped(client, LocalActivityWorkflowImpl.class, "6002");
Thread.sleep(5_000);
System.out.println();
System.out.println("비교: temporal workflow describe -w order-6001 (Regular)");
System.out.println(" temporal workflow describe -w order-6002 (Local)");
System.out.println(" History Length 36 vs 16 / History Size 8712 vs 3104");
}
case "--race" -> {
// [4-8] 스레드 안전하지 않은 액티비티 — 20건 동시 투입
worker.registerWorkflowImplementationTypes(NonIdempotentWorkflowImpl.class);
worker.registerActivitiesImplementations(new UnsafePaymentActivityImpl());
factory.start();
for (int i = 0; i < 20; i++) {
start(client, "7" + String.format("%03d", i), "동시 투입 " + i);
}
System.out.println("콘솔에 '★★★ 오염 감지!' 로그가 뜨는지 보세요.");
}
default -> {
// [4-1] ~ [4-2] 기본 실습
worker.registerWorkflowImplementationTypes(IdempotentWorkflowImpl.class);
worker.registerActivitiesImplementations(new IdempotentPaymentActivityImpl());
factory.start();
start(client, "1001", "기본 실습");
System.out.println("확인: temporal workflow show -w order-1001 --output json \\");
System.out.println(" | jq '.events[4].activityTaskScheduledEventAttributes'");
}
}
}
private static void start(WorkflowClient client, String orderId, String desc) {
OrderWorkflow stub = client.newWorkflowStub(
OrderWorkflow.class,
WorkflowOptions.newBuilder()
.setTaskQueue(TASK_QUEUE)
.setWorkflowId("order-" + orderId)
.build());
OrderRequest req = new OrderRequest(orderId, "C-77", "SKU-A", 2, 39000, " 서울시 강남구 ");
WorkflowClient.start(stub::processOrder, req);
System.out.println("started order-" + orderId + " — " + desc);
}
private static void startTyped(WorkflowClient client, Class<?> type, String orderId) {
OrderWorkflow stub = client.newWorkflowStub(
OrderWorkflow.class,
WorkflowOptions.newBuilder()
.setTaskQueue(TASK_QUEUE)
.setWorkflowId("order-" + orderId)
.build());
OrderRequest req = new OrderRequest(orderId, "C-77", "SKU-A", 2, 39000, " 서울시 강남구 ");
WorkflowClient.start(stub::processOrder, req);
System.out.println("started order-" + orderId + " — " + type.getSimpleName());
}
}
6문제의 문제지입니다. // TODO: 여기에 작성 자리를 채우는 구조입니다.
Ex2NoTimeoutWorkflowImpl 은 그대로 실행하면 워크플로우가 FAILED 로 끝납니다. 먼저 고치기 전에 한 번 돌려서 예외 메시지를 직접 보고 주석에 적은 뒤 고치세요.Ex3BatchActivityImpl 에 heartbeat 와 재개 로직을 넣는 문제입니다. getHeartbeatDetails 를 호출하는 위치가 중요합니다 — 반복문 안에 넣으면 매번 서버를 때립니다.Ex4PaymentActivityImpl 은 chargedOrders 라는 Set 으로 중복을 막으려 합니다. 이 접근이 왜 실패하는지(Worker 가 여러 대이고 재시작되면 Set 이 비어 있습니다) 설명하고, 올바른 방법으로 고치세요.Ex6UnsafeActivityImpl 에는 스레드 안전 문제가 세 군데 있습니다. 하나는 눈에 잘 띄고, 하나는 SimpleDateFormat 이고, 하나는 지연 초기화입니다.package com.example.order;
// =====================================================================================
// Step 04 — 액티비티 : Exercise (6문제)
//
// 실행 방법
// ./gradlew run -PmainClass=com.example.order.Exercise
//
// 규칙
// - `// TODO: 여기에 작성` 자리를 채우세요.
// - 문제 1 과 문제 5 는 코드를 거의 쓰지 않습니다. 표와 근거를 주석으로 채우세요.
// - 정답은 Solution.java 에 있습니다.
// =====================================================================================
import io.temporal.activity.Activity;
import io.temporal.activity.ActivityExecutionContext;
import io.temporal.activity.ActivityInterface;
import io.temporal.activity.ActivityMethod;
import io.temporal.activity.ActivityOptions;
import io.temporal.workflow.Workflow;
import io.temporal.workflow.WorkflowInterface;
import io.temporal.workflow.WorkflowMethod;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.net.http.HttpClient;
import java.text.SimpleDateFormat;
import java.time.Duration;
import java.util.Date;
import java.util.HashSet;
import java.util.Set;
public class Exercise {
public static final String TASK_QUEUE = "ORDER_TASK_QUEUE";
public record OrderRequest(
String orderId, String customerId, String sku, int qty, long amount, String address) {}
@WorkflowInterface
public interface OrderWorkflow {
@WorkflowMethod String processOrder(OrderRequest req);
}
@ActivityInterface
public interface PaymentActivity {
@ActivityMethod(name = "Charge") String charge(String orderId, long amount);
}
@ActivityInterface
public interface BatchActivity {
@ActivityMethod(name = "ProcessBatch") String processBatch(String batchId, int totalCount);
}
// =================================================================================
// 문제 1 — 네 가지 장애 시나리오에 필요한 타임아웃을 고르고 값을 정하세요.
//
// 각 시나리오에 대해
// (1) 어떤 타임아웃이 이 장애를 잡아내는가?
// (2) 값을 몇으로 줄 것인가?
// (3) 왜 다른 타임아웃으로는 못 잡는가?
// 를 주석으로 답하세요.
//
// ┌─────────────────────────────────────────────────────────────────────────────┐
// │ 시나리오 A — 배포 사고로 Worker 가 전부 죽었다. │
// │ 액티비티 태스크가 Task Queue 에 계속 쌓인다. 워크플로우는 전부 Running. │
// │ 에러도 알림도 없다. 고객이 "결제가 안 됐다"고 문의해서야 알게 됐다. │
// │ │
// │ 타임아웃: ______ 값: ______ │
// │ TODO: 여기에 작성 │
// └─────────────────────────────────────────────────────────────────────────────┘
//
// ┌─────────────────────────────────────────────────────────────────────────────┐
// │ 시나리오 B — 외부 결제 API 가 응답을 안 준다. │
// │ 평소 p99 는 1.2초인데, 장애 때는 소켓이 열린 채 영원히 매달린다. │
// │ │
// │ 타임아웃: ______ 값: ______ │
// │ TODO: 여기에 작성 │
// └─────────────────────────────────────────────────────────────────────────────┘
//
// ┌─────────────────────────────────────────────────────────────────────────────┐
// │ 시나리오 C — 2시간짜리 배치 액티비티를 도는 Worker 가 30분 만에 OOM 으로 죽음.│
// │ 서버는 이 사실을 모르고 1시간 30분을 기다린다. │
// │ │
// │ 타임아웃: ______ 값: ______ │
// │ TODO: 여기에 작성 │
// └─────────────────────────────────────────────────────────────────────────────┘
//
// ┌─────────────────────────────────────────────────────────────────────────────┐
// │ 시나리오 D — 결제 액티비티가 재시도를 계속해 하루 종일 돈다. │
// │ 비즈니스 규칙상 "결제는 아무리 재시도해도 30분 안에 끝나거나 실패해야 한다".│
// │ │
// │ 타임아웃: ______ 값: ______ │
// │ TODO: 여기에 작성 │
// └─────────────────────────────────────────────────────────────────────────────┘
//
// 마지막으로, 위 네 가지를 모두 반영한 ActivityOptions 를 아래에 작성하세요.
// =================================================================================
public static ActivityOptions ex1Options() {
// TODO: 여기에 작성
return ActivityOptions.newBuilder()
.build();
}
// =================================================================================
// 문제 2 — 타임아웃을 하나도 안 준 스텁을 고치세요.
//
// 순서
// (1) 먼저 **고치지 말고** 그대로 실행해서 워크플로우가 어떻게 끝나는지 보세요.
// (2) Worker 콘솔의 예외 메시지를 아래 주석에 그대로 적으세요.
// (3) temporal workflow describe -w order-9001 의 Status 를 적으세요.
// (4) 그다음 고치세요.
//
// 예외 메시지: TODO: 여기에 작성
// Status: TODO: 여기에 작성
// 왜 이 에러는 Step 03 의 NonDeterministicException 과 달리
// 워크플로우를 FAILED 로 만드는가? TODO: 여기에 작성
// =================================================================================
public static class Ex2NoTimeoutWorkflowImpl implements OrderWorkflow {
private final PaymentActivity payment = Workflow.newActivityStub(
PaymentActivity.class,
// TODO: 여기에 작성 — 타임아웃을 추가하세요
ActivityOptions.newBuilder().build());
@Override
public String processOrder(OrderRequest req) {
payment.charge(req.orderId(), req.amount());
return req.orderId() + " COMPLETED";
}
}
// =================================================================================
// 문제 3 — 10만 건 배치에 heartbeat 와 재개 로직을 넣으세요.
//
// 조건
// - Worker 가 죽어도 처음부터 다시 하지 않고 중단 지점부터 재개할 것
// - 취소 요청이 오면 감지해서 정리 작업을 하고 종료할 것
// - HeartbeatTimeout 을 얼마로 줄지도 ex3Options() 에서 정할 것
//
// 주의
// getHeartbeatDetails() 를 호출하는 **위치**가 중요합니다.
// 반복문 안에 넣으면 왜 안 되는지 주석으로 적으세요. → TODO: 여기에 작성
// =================================================================================
public static class Ex3BatchActivityImpl implements BatchActivity {
private static final Logger log = LoggerFactory.getLogger(Ex3BatchActivityImpl.class);
@Override
public String processBatch(String batchId, int totalCount) {
ActivityExecutionContext ctx = Activity.getExecutionContext();
// TODO: 여기에 작성 — 이전 시도의 진행 상황을 꺼내세요
int start = 0;
for (int i = start; i < totalCount; i++) {
processOne(batchId, i);
// TODO: 여기에 작성 — 하트비트를 보내고, 취소를 감지하세요
}
return batchId + " DONE " + totalCount;
}
private void processOne(String batchId, int i) {
// 건당 약 2ms 라고 가정
}
private void cleanup(String batchId, int processed) {
log.info("부분 작업 정리 batchId={} processed={}", batchId, processed);
}
}
public static ActivityOptions ex3Options() {
// TODO: 여기에 작성 — StartToClose 와 HeartbeatTimeout 을 정하세요
// 하트비트를 1,000건마다(약 2초마다) 보낸다고 가정합니다.
return ActivityOptions.newBuilder().build();
}
// =================================================================================
// 문제 4 — 멱등하지 않은 결제 액티비티를 고치세요.
//
// 아래 구현은 Set 으로 중복을 막으려 합니다. 이 접근은 실패합니다.
// (a) 왜 실패하는가? 세 가지 이유를 적으세요. → TODO: 여기에 작성
// (b) 멱등키는 어디서 만들어야 하는가? 왜? → TODO: 여기에 작성
// (c) 올바르게 고치세요.
// =================================================================================
public static class Ex4PaymentActivityImpl implements PaymentActivity {
private static final Logger log = LoggerFactory.getLogger(Ex4PaymentActivityImpl.class);
// ✘ 이 방식은 왜 안 되는가?
private final Set<String> chargedOrders = new HashSet<>();
@Override
public String charge(String orderId, long amount) {
if (chargedOrders.contains(orderId)) {
log.info("이미 결제됨 orderId={}", orderId);
return "PAY-DUP";
}
chargedOrders.add(orderId);
log.info("결제 실행 orderId={} amount={}", orderId, amount);
return "PAY-" + orderId;
}
}
// TODO: 여기에 작성 — 올바른 인터페이스와 구현체를 정의하세요.
// 힌트: 인터페이스 시그니처부터 바꿔야 합니다.
//
// @ActivityInterface
// public interface Ex4FixedPaymentActivity { ... }
//
// public static class Ex4FixedPaymentActivityImpl implements Ex4FixedPaymentActivity { ... }
// 그리고 이 워크플로우에서 키를 만들어 넘기세요.
public static class Ex4WorkflowImpl implements OrderWorkflow {
@Override
public String processOrder(OrderRequest req) {
// TODO: 여기에 작성 — 멱등키를 만들고 액티비티에 전달하세요
return req.orderId() + " COMPLETED";
}
}
// =================================================================================
// 문제 5 — 아래 액티비티 5개를 Activity / Local Activity 로 분류하고 근거를 적으세요.
//
// (1) 주소 정규화 — 문자열 trim/대문자 변환. 0.1ms. 외부 호출 없음.
// 분류: ______ 근거: TODO: 여기에 작성
//
// (2) 상품 이미지 리사이징 — 원본 20MB, ImageMagick 호출. 평균 45초. CPU/메모리 많이 씀.
// 분류: ______ 근거: TODO: 여기에 작성
//
// (3) 쿠폰 코드 형식 검증 — 정규식 매칭. 0.05ms. 외부 호출 없음.
// 분류: ______ 근거: TODO: 여기에 작성
//
// (4) 결제 — 외부 PG API 호출. 평균 1.2초. 돈이 움직인다.
// 분류: ______ 근거: TODO: 여기에 작성
//
// (5) 재고 조회 — 내부 재고 서비스 gRPC 호출. 평균 30ms. 읽기 전용.
// 분류: ______ 근거: TODO: 여기에 작성
// ★ 이 항목은 판단이 갈립니다. 어느 쪽을 골랐든 근거를 명확히 적으세요.
//
// 추가로, (2)는 Task Queue 를 분리해야 합니다. 왜인지, 어떻게 하는지 적으세요.
// TODO: 여기에 작성
// =================================================================================
// =================================================================================
// 문제 6 — 스레드 안전 문제를 찾아 고치세요.
//
// 아래 구현체에는 스레드 안전 문제가 **세 군데** 있습니다.
// Worker 는 이 인스턴스 하나를 최대 200개 스레드가 동시에 호출합니다.
// 각 문제 위에 `// [문제 N] 이유` 주석을 달고, 아래에 고친 버전을 작성하세요.
// =================================================================================
public static class Ex6UnsafeActivityImpl implements PaymentActivity {
private static final Logger log = LoggerFactory.getLogger(Ex6UnsafeActivityImpl.class);
// TODO: 여기에 작성 — 문제 지점에 주석을 다세요
private String currentOrderId;
private final SimpleDateFormat fmt = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss");
private HttpClient client;
@Override
public String charge(String orderId, long amount) {
this.currentOrderId = orderId;
if (client == null) {
client = HttpClient.newHttpClient();
}
String at = fmt.format(new Date());
log.info("charge orderId={} amount={} at={}", this.currentOrderId, amount, at);
return "PAY-" + this.currentOrderId;
}
}
// TODO: 여기에 작성 — 고친 버전
//
// public static class Ex6SafeActivityImpl implements PaymentActivity { ... }
public static void main(String[] args) {
System.out.println("Exercise — 각 문제의 TODO 를 채우세요.");
System.out.println("문제 2 는 먼저 고치지 말고 한 번 실행해서 예외를 직접 보세요.");
System.out.println("정답은 Solution.java 를 참고하세요.");
}
}
6문제의 정답과, "왜 그런지"를 설명하는 긴 주석이 들어 있습니다. 문제를 풀어 본 뒤에 여세요.
IllegalStateException: Either ScheduleToCloseTimeout or StartToCloseTimeout is required 입니다. 이 실패가 워크플로우를 FAILED 로 만든다는 점(Step 03 의 NonDeterministicException 과 다름)과, 그래서 로컬에서 한 번만 돌려 봐도 즉시 발견된다는 점을 대비해 설명합니다.getHeartbeatDetails() 를 반복문 밖에서 딱 한 번 호출합니다. 이 값은 액티비티 시도가 시작될 때 서버가 함께 내려 주는 것이라 반복 호출할 이유가 없습니다. 하트비트 주기(1,000건마다)를 정하는 계산 근거도 함께 적었습니다.Set 접근이 왜 틀렸는지를 세 단계로 설명합니다. (1) Worker 가 여러 대면 Set 이 공유되지 않는다 (2) Worker 재시작이면 비어 있다 (3) 애초에 액티비티는 무상태여야 한다. 올바른 답은 DB 유니크 제약 + INSERT ... ON CONFLICT DO NOTHING 이거나 PG 의 Idempotency-Key 이며, 키는 워크플로우가 Workflow.randomUUID() 로 만들어 넘깁니다.private String currentOrderId 요청별 상태 (2) private final SimpleDateFormat fmt — SimpleDateFormat 은 스레드 안전하지 않아 200 스레드가 쓰면 날짜가 뒤섞입니다. DateTimeFormatter 로 교체합니다 (3) if (client == null) client = new HttpClient() 지연 초기화 — 동시 진입 시 클라이언트가 여러 개 만들어져 커넥션이 샙니다. 생성자 주입으로 바꿉니다.package com.example.order;
// =====================================================================================
// Step 04 — 액티비티 : Solution (6문제 정답 + 해설)
//
// 실행 방법
// ./gradlew run -PmainClass=com.example.order.Solution
//
// Exercise.java 를 먼저 풀어 본 뒤 읽으세요.
// =====================================================================================
import io.temporal.activity.Activity;
import io.temporal.activity.ActivityExecutionContext;
import io.temporal.activity.ActivityInterface;
import io.temporal.activity.ActivityMethod;
import io.temporal.activity.ActivityOptions;
import io.temporal.activity.LocalActivityOptions;
import io.temporal.client.ActivityCompletionException;
import io.temporal.common.RetryOptions;
import io.temporal.workflow.Workflow;
import io.temporal.workflow.WorkflowInterface;
import io.temporal.workflow.WorkflowMethod;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.net.http.HttpClient;
import java.time.Duration;
import java.time.Instant;
import java.time.ZoneId;
import java.time.format.DateTimeFormatter;
public class Solution {
public static final String TASK_QUEUE = "ORDER_TASK_QUEUE";
public static final String MEDIA_TASK_QUEUE = "MEDIA_TASK_QUEUE";
public record OrderRequest(
String orderId, String customerId, String sku, int qty, long amount, String address) {}
@WorkflowInterface
public interface OrderWorkflow {
@WorkflowMethod String processOrder(OrderRequest req);
}
// =================================================================================
// 정답 1 — 시나리오별 타임아웃
//
// ┌───────────────────────────────────────────────────────────────────────────────┐
// │ A. Worker 전멸 → **ScheduleToStart = 5분** │
// │ │
// │ 이 절의 핵심입니다. StartToClose 를 아무리 짧게 줘도 이 장애는 못 잡습니다. │
// │ StartToClose 는 "액티비티가 **시작한 뒤부터**" 재는 타이머인데, Worker 가 │
// │ 없으면 애초에 시작을 못 하므로 타이머가 돌지 않습니다. │
// │ │
// │ 실제로 이렇게 됩니다. │
// │ $ temporal task-queue describe --task-queue ORDER_TASK_QUEUE \ │
// │ --task-queue-type activity │
// │ BuildId TaskQueueType Pollers BacklogCount ApproximateBacklogAge │
// │ ACTIVITY 0 84291 3h 42m 11s │
// │ │
// │ 워크플로우는 전부 Running, 액티비티는 3시간 42분째 대기. 에러도 알림도 없음. │
// │ ScheduleToStart 5분이면 5분 만에 ActivityTaskTimedOut 이 터지고, 워크플로우 │
// │ 실패율 알림이 즉시 울립니다. │
// │ │
// │ ★ ScheduleToStart 는 성능 옵션이 아니라 **관측 장치**입니다. │
// │ ★ 이 타임아웃은 재시도되지 않습니다(RETRY_STATE_NON_RETRYABLE_FAILURE). │
// │ 재시도해도 어차피 같은 빈 큐에 다시 들어갈 뿐이기 때문입니다. │
// │ ★ 오타 난 Task Queue 이름도 이걸로 잡힙니다. 가장 흔한 신규 개발자 실수입니다.│
// └───────────────────────────────────────────────────────────────────────────────┘
//
// ┌───────────────────────────────────────────────────────────────────────────────┐
// │ B. 외부 API 무응답 → **StartToClose = 5초** (p99 1.2초 × 4배) │
// │ │
// │ StartToClose 는 "한 번의 시도"를 재므로 매달린 시도를 끊고 재시도합니다. │
// │ 값은 p99 × 2~3배가 기본입니다. 여기서는 여유를 둬 4배. │
// │ 너무 짧으면 정상 요청까지 죽여 재시도 폭풍이 나고, 너무 길면 장애 감지가 │
// │ 늦습니다. │
// │ ScheduleToStart 로는 못 잡습니다 — 이미 시작한 액티비티이기 때문입니다. │
// └───────────────────────────────────────────────────────────────────────────────┘
//
// ┌───────────────────────────────────────────────────────────────────────────────┐
// │ C. 배치 Worker OOM → **HeartbeatTimeout = 30초** │
// │ │
// │ StartToClose 2시간은 유지합니다(정상 실행에 필요하므로). 하트비트가 있으면 │
// │ 30초 만에 죽음을 감지합니다. **1시간 30분 → 30초.** │
// │ 덤으로 lastHeartbeatDetails 가 남아 재시도가 중단 지점부터 이어집니다. │
// │ │
// │ ★ 값은 하트비트 주기의 2~3배로 잡습니다. 30초마다 뛰면 60~90초. │
// │ ★ HeartbeatTimeout 을 설정하지 않으면 heartbeat() 호출 자체가 무효입니다. │
// └───────────────────────────────────────────────────────────────────────────────┘
//
// ┌───────────────────────────────────────────────────────────────────────────────┐
// │ D. 재시도 총량 제한 → **ScheduleToClose = 30분** │
// │ │
// │ ScheduleToClose 만이 "재시도를 전부 포함한 총 시간"을 잽니다. │
// │ RetryOptions.setMaximumAttempts() 로도 제한할 수 있지만, 백오프 간격에 따라 │
// │ 총 시간이 들쭉날쭉합니다. **비즈니스 데드라인은 시간으로 거세요.** │
// │ │
// │ 터졌을 때 히스토리의 failure.cause 에 직전 시도의 실패 이유가 중첩됩니다. │
// │ "message": "activity ScheduleToClose timeout" │
// │ "cause": { "message": "activity StartToClose timeout" } │
// └───────────────────────────────────────────────────────────────────────────────┘
// =================================================================================
public static ActivityOptions sol1Options() {
return ActivityOptions.newBuilder()
.setScheduleToStartTimeout(Duration.ofMinutes(5)) // A
.setStartToCloseTimeout(Duration.ofSeconds(5)) // B
.setHeartbeatTimeout(Duration.ofSeconds(30)) // C (장기 액티비티에)
.setScheduleToCloseTimeout(Duration.ofMinutes(30)) // D
.build();
}
// =================================================================================
// 정답 2 — 타임아웃 필수 규칙
//
// 예외 메시지
// io.temporal.failure.ApplicationFailure: message='Either ScheduleToCloseTimeout
// or StartToCloseTimeout is required', type='java.lang.IllegalStateException'
// Caused by: java.lang.IllegalStateException: Either ScheduleToCloseTimeout or
// StartToCloseTimeout is required
// at io.temporal.common.interceptors.ActivityOptionsUtils
// .validateAndBuildOptions(ActivityOptionsUtils.java:61)
//
// Status: FAILED
//
// 왜 FAILED 가 되는가 (Step 03 의 NonDeterministicException 과의 대비)
// NonDeterministicException 은 SDK 가 "코드와 히스토리가 안 맞는다"고 판단한
// **인프라 수준의 문제**입니다. 코드를 고쳐 재배포하면 복구되므로 워크플로우를
// 죽이지 않고 Workflow Task 만 무한 재시도합니다(Status: RUNNING).
//
// 반면 이 예외는 **워크플로우 코드가 던진 애플리케이션 예외**로 취급됩니다.
// SDK 입장에서는 "개발자가 잘못된 옵션을 넘겼다"이고, 재시도해도 결과가 같습니다.
// 그래서 WorkflowExecutionFailed 로 닫습니다.
//
// 실무적 의미
// 이 실수는 **로컬에서 한 번만 돌려 봐도 즉시 발견됩니다.** 워크플로우가 바로
// FAILED 로 끝나기 때문입니다. Step 03 의 결정성 위반이 무서운 이유는 정확히
// 그 반대 — 로컬에서 100% 통과하고 운영 재배포 때만 터지기 때문입니다.
// =================================================================================
@ActivityInterface
public interface PaymentActivity {
@ActivityMethod(name = "Charge")
String charge(String orderId, long amount);
}
public static class Sol2FixedWorkflowImpl implements OrderWorkflow {
private final PaymentActivity payment = Workflow.newActivityStub(
PaymentActivity.class,
ActivityOptions.newBuilder()
// 최소한 이 둘 중 하나. 실무에서는 셋 다 주는 것이 원칙.
.setStartToCloseTimeout(Duration.ofSeconds(10))
.setScheduleToStartTimeout(Duration.ofMinutes(5))
.setScheduleToCloseTimeout(Duration.ofMinutes(30))
.build());
@Override
public String processOrder(OrderRequest req) {
payment.charge(req.orderId(), req.amount());
return req.orderId() + " COMPLETED";
}
}
// =================================================================================
// 정답 3 — heartbeat 와 재개
//
// getHeartbeatDetails() 를 **반복문 밖에서 딱 한 번** 호출합니다.
//
// 왜 반복문 안에 넣으면 안 되는가
// (1) 이 값은 액티비티 시도가 시작될 때 서버가 ActivityTask 에 실어 함께
// 내려 준 것입니다. 시도가 도는 동안 바뀌지 않습니다. 반복 호출은
// 같은 값을 다시 읽는 낭비일 뿐입니다.
// (2) 더 나쁜 것은 로직 오류입니다. 반복문 안에서 start 를 다시 읽어 i 에
// 대입하면 진행이 되돌아가 무한 루프가 됩니다.
//
// 하트비트 주기 계산
// 건당 2ms × 1,000건 = 약 2초마다 heartbeat() 호출.
// SDK 는 HeartbeatTimeout 의 80% 간격으로 묶어 실제 전송하므로,
// HeartbeatTimeout 60초면 실제 gRPC 는 48초마다 한 번입니다.
// 10만 건 = 100회 호출 → 실제 전송은 4~5회. 서버 부하가 아닙니다.
//
// HeartbeatTimeout 은 하트비트 주기(2초)의 2~3배가 아니라, **여유롭게** 잡습니다.
// GC 정지나 일시적 지연으로 하트비트가 늦을 수 있기 때문입니다. 60초를 권장합니다.
//
// StartToClose 계산
// 10만 건 × 2ms = 200초. 여유를 둬 10분.
// ★ 하트비트가 있으면 StartToClose 를 넉넉히 줘도 안전합니다. 죽은 Worker 는
// 하트비트가 잡아 주기 때문입니다. 하트비트가 없다면 StartToClose 를 짧게
// 줘야 하는데, 그러면 정상 실행까지 죽이게 됩니다. 이 딜레마를 하트비트가 풉니다.
// =================================================================================
@ActivityInterface
public interface BatchActivity {
@ActivityMethod(name = "ProcessBatch")
String processBatch(String batchId, int totalCount);
}
public static class Sol3BatchActivityImpl implements BatchActivity {
private static final Logger log = LoggerFactory.getLogger(Sol3BatchActivityImpl.class);
@Override
public String processBatch(String batchId, int totalCount) {
ActivityExecutionContext ctx = Activity.getExecutionContext();
// ★ 반복문 밖에서 한 번만.
int start = ctx.getHeartbeatDetails(Integer.class).orElse(0);
int attempt = ctx.getInfo().getAttempt();
log.info("배치 시작 batchId={} start={} total={} attempt={}",
batchId, start, totalCount, attempt);
for (int i = start; i < totalCount; i++) {
processOne(batchId, i);
if (i % 1000 == 0) {
try {
ctx.heartbeat(i);
} catch (ActivityCompletionException e) {
// 취소 감지. ActivityCanceledException 이 여기로 온다.
// 취소 전달 경로는 하트비트뿐이라, 하트비트가 없는 액티비티는
// 취소 요청이 와도 끝까지 실행된다.
log.info("취소 감지 — {}건 처리 후 정리하고 종료", i);
cleanup(batchId, i);
throw e; // 반드시 다시 던져야 서버가 ActivityTaskCanceled 로 기록한다
}
}
}
return batchId + " DONE " + totalCount;
}
private void processOne(String batchId, int i) {}
private void cleanup(String batchId, int processed) {
log.info("부분 작업 정리 batchId={} processed={}", batchId, processed);
}
}
public static ActivityOptions sol3Options() {
return ActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofMinutes(10)) // 10만 건 × 2ms = 200초 + 여유
.setHeartbeatTimeout(Duration.ofSeconds(60)) // 2초마다 뛰는데 60초 여유
.setScheduleToStartTimeout(Duration.ofMinutes(5))
.build();
}
// =================================================================================
// 정답 4 — 멱등성
//
// (a) Set 접근이 실패하는 세 가지 이유
//
// 1. **Worker 가 여러 대면 Set 이 공유되지 않는다.**
// 첫 시도를 Worker A 가, 재시도를 Worker B 가 집어 갑니다. B 의 Set 은 비어
// 있으므로 그대로 결제합니다. 운영은 Worker 를 항상 여러 대 띄웁니다.
//
// 2. **Worker 재시작이면 비어 있다.**
// Set 은 프로세스 메모리입니다. 배포·OOM·스케일인으로 프로세스가 바뀌면
// 전부 사라집니다. 그리고 재시도가 필요한 상황은 대개 Worker 가 죽은 상황입니다.
// 즉 **가장 필요한 순간에 반드시 비어 있습니다.**
//
// 3. **애초에 액티비티는 무상태여야 한다.**
// Worker 는 액티비티 인스턴스 하나를 최대 200 스레드가 공유합니다.
// HashSet 은 스레드 안전하지도 않아 동시 쓰기에 데이터가 깨집니다.
//
// 추가로, orderId 를 키로 쓴 것도 문제입니다. 같은 주문에 대해 부분 환불 후
// 재결제 같은 정당한 재시도가 필요할 수 있는데, 그것까지 막아 버립니다.
// 키는 **"이 액티비티 호출"** 단위여야지 **"이 주문"** 단위가 아닙니다.
//
// (b) 멱등키는 **워크플로우**에서 만듭니다.
//
// 이유는 히스토리에 있습니다. 워크플로우가 만들어 액티비티 인자로 넘기면,
// 그 값이 ActivityTaskScheduled 이벤트의 input 에 기록됩니다.
//
// "input": { "payloads": [
// { "data": "IjEwMDEi" }, ← orderId
// { "data": "MzkwMDA=" }, ← amount
// { "data": "IjNmMmE5YzE0LTdiMDEtNDhkNC1hZWY2LTkwMGMxZDJmODhhNCI=" } ← 멱등키
// ] }
//
// ActivityTaskScheduled 는 재시도할 때 **새로 생기지 않습니다.** 서버가 같은
// 이벤트의 input 으로 다시 디스패치할 뿐이라, 몇 번을 재시도해도 액티비티는
// 같은 키를 받습니다.
//
// 액티비티 안에서 UUID 를 만들면 시도마다 새 키가 되어 아무 의미가 없습니다.
//
// 그리고 반드시 **Workflow.randomUUID()** 여야 합니다.
// UUID.randomUUID() 는 Worker 재시작으로 리플레이가 일어날 때 값이 바뀝니다.
// (Step 03 의 3-7 함정 참조)
//
// (c) 올바른 구현 — 아래 두 가지 중 하나
// 1. 외부 API 의 Idempotency-Key 헤더 (PG, 이메일 발송 등 대부분 지원)
// 2. DB 유니크 제약 + INSERT ... ON CONFLICT DO NOTHING (내부 시스템)
// =================================================================================
@ActivityInterface
public interface Sol4PaymentActivity {
@ActivityMethod(name = "ChargeIdempotent")
String chargeIdempotent(String orderId, long amount, String idempotencyKey);
}
public static class Sol4PaymentActivityImpl implements Sol4PaymentActivity {
private static final Logger log = LoggerFactory.getLogger(Sol4PaymentActivityImpl.class);
// 스레드 안전한 불변 의존성만 필드에 둔다.
private final HttpClient http;
public Sol4PaymentActivityImpl(HttpClient http) {
this.http = http;
}
@Override
public String chargeIdempotent(String orderId, long amount, String idempotencyKey) {
log.info("charge orderId={} amount={} key={}", orderId, amount, idempotencyKey);
// 방법 1 — PG 의 Idempotency-Key 헤더
// HttpRequest req = HttpRequest.newBuilder()
// .uri(URI.create("https://pg.example.com/v1/charges"))
// .header("Idempotency-Key", idempotencyKey)
// .POST(...)
// .build();
// 같은 키가 두 번 오면 PG 가 저장된 첫 응답을 그대로 돌려준다. 돈은 한 번만.
// 방법 2 — 내부 DB 라면 유니크 제약
// INSERT INTO payments (idempotency_key, order_id, amount, status)
// VALUES (?, ?, ?, 'CHARGED')
// ON CONFLICT (idempotency_key) DO NOTHING
// RETURNING payment_id;
// → 0행이 돌아오면 이미 처리된 것. 기존 payment_id 를 조회해 반환한다.
return "PAY-" + idempotencyKey.substring(0, 8);
}
}
public static class Sol4WorkflowImpl implements OrderWorkflow {
private static final Logger log = Workflow.getLogger(Sol4WorkflowImpl.class);
private final Sol4PaymentActivity payment = Workflow.newActivityStub(
Sol4PaymentActivity.class,
ActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofSeconds(10))
.setScheduleToCloseTimeout(Duration.ofMinutes(30))
.setRetryOptions(RetryOptions.newBuilder().setMaximumAttempts(5).build())
.build());
@Override
public String processOrder(OrderRequest req) {
// ★ Workflow.randomUUID() — Run ID 시드. 리플레이해도 같은 값.
String idemKey = Workflow.randomUUID().toString();
log.info("[{}] 멱등키 발급 {}", req.orderId(), idemKey);
String paymentId = payment.chargeIdempotent(req.orderId(), req.amount(), idemKey);
log.info("[{}] 결제 완료 {}", req.orderId(), paymentId);
return req.orderId() + " COMPLETED";
}
}
// =================================================================================
// 정답 5 — Activity / Local Activity 분류
//
// (1) 주소 정규화 → **Local Activity**
// 순수 문자열 처리, 0.1ms, 외부 호출 없음. 재시도해도 부담 없음.
// MarkerRecorded 1개로 끝나므로 히스토리 비용이 1/3 이하.
//
// (2) 이미지 리사이징 → **Activity** (그리고 Task Queue 분리)
// 45초는 Workflow Task Timeout(기본 10초)을 훨씬 넘습니다. Local Activity 로
// 하면 SDK 가 Workflow Task 를 계속 연장하려다 히스토리만 지저분해지고,
// Worker 가 죽으면 진행 상황이 통째로 사라집니다.
// Activity 로 하면 하트비트로 진행 상황을 저장하고, 죽어도 이어서 할 수 있습니다.
//
// Task Queue 분리가 필요한 이유
// 이미지 변환은 CPU·메모리를 크게 씁니다. 결제 액티비티와 같은 Worker 에서
// 돌면 배치 하나가 결제 전체를 굶깁니다(maxConcurrentActivityExecutionSize 를
// 점유). Worker 가 OOM 으로 죽으면 그 위의 결제 액티비티도 함께 죽습니다.
//
// // 워크플로우 코드에서 액티비티별로 다른 Task Queue 를 지정
// Workflow.newActivityStub(MediaActivity.class,
// ActivityOptions.newBuilder()
// .setTaskQueue("MEDIA_TASK_QUEUE")
// .setStartToCloseTimeout(Duration.ofMinutes(30))
// .setHeartbeatTimeout(Duration.ofSeconds(30))
// .build());
//
// 그리고 고사양 인스턴스에 MEDIA_TASK_QUEUE 만 폴링하는 Worker 를 따로 띄웁니다.
// 워크플로우 코드는 한 줄도 안 바뀝니다.
//
// (3) 쿠폰 코드 형식 검증 → **Local Activity**
// (1)과 같은 이유. 사실 이 정도면 **워크플로우 코드에 직접 써도** 됩니다.
// 정규식 매칭은 결정적이므로 액티비티로 뺄 이유가 없습니다.
// 액티비티로 빼는 것은 "결정적이지 않거나, 실패할 수 있거나, 느릴 때"입니다.
//
// (4) 결제 → **Activity**
// 외부 API + 돈이 움직임. 재시도·타임아웃·가시성이 전부 필요합니다.
// Local Activity 는 재시도가 Worker 메모리에서 관리되어 히스토리에 안 남고,
// Worker 가 죽으면 재시도 상태가 사라집니다. 돈이 걸린 일에 쓸 수 없습니다.
//
// (5) 재고 조회 → **Activity 권장** (판단이 갈리는 항목)
//
// Local Activity 를 고른 근거도 타당합니다: 읽기 전용, 30ms, 재시도 부담 없음.
// 그러나 **일반 Activity 를 권합니다.**
//
// - 외부 시스템(재고 서비스)을 건드립니다. 그 서비스가 장애일 때 재시도
// 횟수·간격·타임아웃이 히스토리에 남아야 원인 분석이 됩니다.
// Local Activity 는 재시도가 Worker 메모리에서만 일어나 흔적이 없습니다.
// - 30ms 는 정상일 때의 값입니다. 장애 시에는 30초가 됩니다. 그러면
// Workflow Task Timeout 을 넘겨 문제가 됩니다.
// - 재고 서비스가 느려질 때 Task Queue 를 분리해 격리할 여지가 필요합니다.
//
// 기준을 한 문장으로: **"외부 시스템을 건드리면 일반 Activity."**
// Local Activity 는 "이 Worker 프로세스 안에서 완결되는 순수 계산"에만 씁니다.
// =================================================================================
@ActivityInterface
public interface Sol5ValidationActivity {
@ActivityMethod(name = "NormalizeAddress") String normalizeAddress(String address);
@ActivityMethod(name = "ValidateCoupon") boolean validateCoupon(String code);
}
public static class Sol5WorkflowImpl implements OrderWorkflow {
// (1)(3) 순수 계산 → Local Activity
private final Sol5ValidationActivity validation = Workflow.newLocalActivityStub(
Sol5ValidationActivity.class,
LocalActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofSeconds(5))
.build());
// (4) 결제 → 일반 Activity
private final Sol4PaymentActivity payment = Workflow.newActivityStub(
Sol4PaymentActivity.class,
ActivityOptions.newBuilder()
.setStartToCloseTimeout(Duration.ofSeconds(10))
.setScheduleToStartTimeout(Duration.ofMinutes(5))
.setScheduleToCloseTimeout(Duration.ofMinutes(30))
.build());
@Override
public String processOrder(OrderRequest req) {
String addr = validation.normalizeAddress(req.address());
String key = Workflow.randomUUID().toString();
payment.chargeIdempotent(req.orderId(), req.amount(), key);
return req.orderId() + " COMPLETED " + addr;
}
}
// =================================================================================
// 정답 6 — 스레드 안전 문제 세 군데
//
// [문제 1] private String currentOrderId;
// 요청별 상태를 인스턴스 필드에 담았습니다. Worker 는 이 인스턴스 하나를
// 최대 200 스레드가 동시에 호출합니다.
//
// 스레드 A: currentOrderId = "1001"
// 스레드 B: currentOrderId = "1002" ← 덮어씀
// 스레드 A: return "PAY-" + currentOrderId → "PAY-1002" ★ 남의 주문
//
// **다른 고객의 주문 ID 로 결제 결과가 나갑니다.** 재현이 어렵고, 부하가
// 높을 때만 나타나며, 로그만 봐서는 원인을 찾기 힘든 최악의 버그입니다.
// → 지역 변수와 파라미터만 씁니다.
//
// [문제 2] private final SimpleDateFormat fmt = new SimpleDateFormat(...);
// SimpleDateFormat 은 **스레드 안전하지 않습니다.** 내부에 Calendar 를 필드로
// 들고 있어 동시 호출 시 서로의 중간 상태를 읽습니다. 결과는
// - 날짜가 뒤섞인 문자열 ("2026-03-11 09:2026")
// - NumberFormatException / ArrayIndexOutOfBoundsException 이 무작위로
// final 이라 안전해 보이는 것이 함정입니다. **참조가 불변일 뿐 객체는 가변**입니다.
// → java.time.format.DateTimeFormatter 는 불변이고 스레드 안전합니다.
// ZoneId 도 명시하세요. 안 하면 컨테이너 TZ 에 따라 결과가 달라집니다.
//
// [문제 3] if (client == null) { client = HttpClient.newHttpClient(); }
// 지연 초기화(lazy init)에 동기화가 없습니다. 두 스레드가 동시에 null 체크를
// 통과하면 HttpClient 가 두 개 만들어지고, 하나는 참조를 잃은 채 커넥션 풀과
// 셀렉터 스레드를 그대로 들고 있습니다. 요청이 몰릴수록 이 누수가 쌓입니다.
// 더 나쁜 것은 client 가 volatile 도 아니라, 다른 스레드가 **부분적으로 초기화된
// 객체**를 볼 수 있다는 점입니다.
// → 생성자 주입으로 바꿉니다. 액티비티 인스턴스는 Worker 등록 시 한 번만
// 만들어지므로 지연 초기화를 할 이유가 없습니다.
//
// 정리: 액티비티 구현체의 필드에는 **불변이거나 스레드 안전한 것만** 둡니다.
// HttpClient, DateTimeFormatter, Spring Repository, DataSource 는 안전.
// SimpleDateFormat, StringBuilder, HashMap, 요청별 상태는 위험.
// =================================================================================
public static class Sol6SafeActivityImpl implements PaymentActivity {
private static final Logger log = LoggerFactory.getLogger(Sol6SafeActivityImpl.class);
// 불변 + 스레드 안전한 의존성만. 전부 생성자 주입.
private final HttpClient http;
private final DateTimeFormatter fmt;
public Sol6SafeActivityImpl(HttpClient http) {
this.http = http;
this.fmt = DateTimeFormatter.ofPattern("yyyy-MM-dd HH:mm:ss")
.withZone(ZoneId.of("Asia/Seoul"));
}
@Override
public String charge(String orderId, long amount) {
// 모든 상태는 지역 변수로
String at = fmt.format(Instant.now());
log.info("charge orderId={} amount={} at={}", orderId, amount, at);
return "PAY-" + orderId;
}
}
public static void main(String[] args) {
System.out.println("Solution — 각 정답 클래스를 Worker 에 등록해 실행하세요.");
System.out.println("정답 1 과 정답 5 의 해설은 주석에 있습니다.");
}
}