주요 컨텐츠로 이동
파트너

Temporal 및 Lakebase를 사용하여 내구성 있는 에이전트 구축하기

작성자: Sam Ingbar

  • Temporal로 에이전트 진행 상황 보존: 워커 실패 후 기록된 작업을 복구하고, 실패한 작업을 재시도하며, 사람의 검토를 안정적으로 대기합니다.
  • Lakebase Postgres로 실시간 애플리케이션 상태 제공: 각 실행 전반에 걸쳐 근거, 추천, 검토 결정 및 운영 메트릭을 쿼리할 수 있도록 합니다.
  • 실행을 거버넌스가 적용된 데이터에 연결: 동기화된 테이블을 통해 Unity Catalog 정책을 읽고, Lakebase Change Data Feed를 활성화하여 운영 변경 사항을 Delta 이력 테이블에 다시 게시합니다.

개인 대출 심사 에이전트는 증거를 수집하고, 정책을 적용하며, 검토자를 며칠 동안 기다릴 수도 있습니다. 그 동안 워커(worker)가 재시작되거나 도구 호출이 실패할 수 있습니다. 애플리케이션은 완료된 작업을 보존하고, 실행을 재개하며, 검토자가 증거를 계속 사용할 수 있도록 유지해야 합니다.

참조 구현은 지속 가능한 실행을 위해 Temporal을 사용하고, 쿼리 가능한 운영 상태를 위해 Lakebase Postgres를 사용합니다. 동기화된 테이블을 통해 Unity Catalog의 심사 정책을 Lakebase에서 사용할 수 있습니다. Temporal Activities는 증거, 결정 및 메트릭을 Lakebase에 기록합니다. 활성화되면 Lakebase Change Data Feed가 이러한 변경 사항을 Unity Catalog에서 관리하는 Delta 기록 테이블에 게시할 수 있습니다. 이 조합은 Databricks가 에이전트의 입력과 다운스트림 분석을 이미 관리하고 있을 때 특히 유용합니다.

장시간 실행되는 클라우드 에이전트의 과제

클라우드 에이전트는 자신을 시작한 요청, 워커, 컨테이너 또는 배포보다 더 오래 유지될 수 있습니다. 사용자는 세션을 시작하고 다음 날 돌아와 다른 워커에서 계속 진행할 수 있습니다. 배포 및 프로세스 실패는 일상적인 일이므로, 에이전트의 진행 상황은 이를 실행하는 프로세스와 독립적으로 유지되어야 합니다. 복구하려면 완료된 작업의 결과와 다음에 수행할 작업을 결정하는 데 필요한 제어 흐름 상태가 모두 필요합니다.

이 심사 에이전트의 경우, 다음과 같은 6가지 요구사항이 발생합니다.

  1. 복구: 대체 워커가 마지막으로 완료된 단계부터 재개해야 합니다.
  2. 재시도: 도구 호출 및 데이터베이스 작업은 부작용을 중복해서 일으키지 않고 반복 실행을 견뎌낼 수 있어야 합니다.
  3. 장시간 대기: 에이전트는 워커를 계속 점유하지 않고 사람이나 외부 시스템을 기다릴 수 있어야 합니다.
  4. 운영 가시성: 애플리케이션과 운영자는 현재 상태, 증거, 재시도 상태 및 실패 세부 정보가 필요합니다.
  5. 런타임 거버넌스: 코드 배포 없이도 정책 업데이트를 적용할 수 있어야 하며, 애플리케이션은 진행 중인 케이스가 이를 언제 채택할지 정의해야 합니다.
  6. 감사: 시스템은 각 실행과 관련된 증거, 정책, 권장 사항 및 사람의 결정을 보존해야 합니다.

대화 기록은 이러한 상태의 일부만 다룹니다. 복구하려면 제어 흐름 기록도 필요합니다. 즉, 어떤 작업이 예약되었는지, 어떤 결과가 기록되었는지, 에이전트가 무엇을 기다리고 있는지, 어떤 명령을 수락했는지 등이 포함됩니다.

Temporal은 분산 시스템 관리를 단순화합니다. Temporal로 빌드할 때 Workflow는 단일 에이전트 실행을 위한 지속 가능한 제어 흐름입니다. Activity는 모델, 도구 또는 데이터베이스에 대한 호출로, 그 결과는 Workflow의 Event History에 기록되며 Activity는 재시도할 수 있습니다. Signal은 심사역의 결정과 같이 실행 중인 Workflow에 전송되는 비동기 명령입니다. Temporal Lakebase AgentWorkflow 참조 구현은 실행 가능한 개인 대출 심사 에이전트입니다. 이 에이전트는 여러 도구를 호출하고, 거버넌스가 적용된 정책을 읽고, 권장 사항을 생성하고, 심사역을 기다립니다.

Lakebase Postgres도 개발자가 이러한 문제를 관리하는 데 도움을 주지만, Temporal과 Lakebase는 서로 다른 소비자를 위해 서로 다른 상태를 저장합니다. Temporal의 Event History는 재실행(replay)을 구동합니다. Lakebase는 애플리케이션 지향 뷰(현재 실행 상태, 메시지, 증거, 검토 상태 및 메트릭)를 저장합니다. Unity Catalog는 정책 소스로 유지되며, 동기화된 테이블을 통해 해당 정책을 Postgres에서 쿼리할 수 있고, Change Data Feed는 운영 기록의 반환 경로를 제공합니다. 두 시스템은 트랜잭션을 공유하지 않습니다. Lakebase 쓰기는 최소 한 번 실행(at-least-once execution) 방식으로 Temporal Activities로 실행됩니다. 결정론적 식별자, 제약 조건, 보호된 업데이트 및 Postgres upsert를 통해 반복되는 Activity 시도가 동일한 논리적 레코드를 대상으로 하도록 보장합니다.

이 아키텍처는 두 개의 관리형 시스템과 이들 간의 프로젝션 계약을 추가합니다. 이들은 함께 운영 오버헤드를 낮게 유지하면서 에이전트의 복원력과 확장성을 향상시킵니다. Temporal과 Lakebase의 조합은 에이전트 세션이 워커 교체 후에도 유지되어야 하고, 오랜 대기 후에 입력을 수락해야 하며, 관계형 상태를 애플리케이션에 노출하고, 열려 있는 동안 거버넌스가 적용된 데이터를 적용해야 할 때 가장 유용합니다.

대출 심사 사용 사례

대출 심사를 선택한 이유는 동일한 실행에서 증거를 수집하고, 정책을 적용하고, 권장 사항을 생성하고, 사람을 기다려야 하기 때문입니다. 이러한 단계 중 어느 단계에서든 워커가 실패할 수 있습니다. 애플리케이션 배포 없이도 정책이 변경될 수 있으며, UI는 Workflow가 종료되기 전에 현재 증거를 필요로 합니다.

모의(mock) 신청자가 실제 신용 조사 기관 및 소득 제공업체를 대체하며, 단순화를 위해 도구 순서는 결정론적입니다. 각 요청에는 사용자 ID, 신청자 ID, 금액, 목적, 모델 선택 및 턴 제한이 포함됩니다. FastAPI는 run_id을(를) 할당하고, LoanUnderwritingWorkflow을(를) 시작하며, API, Temporal 실행 및 Lakebase 행 전체에서 동일한 ID를 사용합니다.

첫 번째 턴에서 credit_check은(는) 점수, 거래 내역, 연체 정보 및 현재 부채를 반환합니다. income_verification은(는) 소득 및 고용 증빙을 반환합니다. debt_to_income_calc은(는) 총부채상환비율(DTI)을 계산합니다. policy_lookup은(는) 대출 목적에 맞는 정책을 로드하고 승인, 회부 및 거절 임계값에 따라 증거를 평가합니다.

샘플의 경계선상에 있는 신청자는 신용 점수 665점, 검증된 연간 소득 $76,000, 월 부채 $2,400 및 중요하지 않은 연체 플래그 1개를 가지고 있습니다. 정책 결과에는 모든 규칙, 임계값, 실제 값, 합격/불합격 결과, 출처, 권장 사항 및 근거가 기록됩니다. 모델은 권장할 수는 있지만 결정할 수는 없습니다. 심사역은 승인, 거부 또는 추가 정보를 요청합니다. 추가 정보 요청은 또 다른 사용자 메시지와 또 다른 에이전트 턴이 됩니다. 이 케이스는 완료된 도구 호출 후 워커 충돌, Activity 완료가 손실된 커밋된 Lakebase 쓰기, 며칠 동안 열려 있는 검토, 만료된 브라우저 결정, 실행 중 정책 변경 등을 테스트합니다.

아키텍처

image1.jpg
그림 1. 참조 구현의 실행, 운영 상태 및 거버넌스 경로.

심사 에이전트를 구현하기 위해 React와 FastAPI가 HTTP 및 UI 작업을 처리합니다(실행 시작, 증거 렌더링, 케이스 목록 표시, 검토 결정 제출). Temporal Cloud는 Event History를 저장하고 Task를 발송합니다. 워커는 Workflow 코드를 재실행하고 모델, 도구 및 Lakebase Activities를 실행합니다. 네트워크 및 데이터베이스 I/O는 결정론적 Workflow 코드 외부에 유지됩니다.

실행은 FastAPI가 Workflow를 시작할 때 시작됩니다. 워커가 Activities를 예약하고, Temporal이 그 결과를 기록하며, 에이전트는 결국 AWAITING_REVIEW에 도달합니다. 심사역의 응답은 Signal을 통해 반환됩니다. 승인 또는 거부는 실행을 종료하고, 추가 정보 요청은 에이전트 루프를 재개합니다.

Lakebase는 두 개의 운영 스키마를 보유합니다. agent_ops에는 FastAPI가 SQL로 쿼리할 수 있는 실행 상태, 메시지, 도구 호출, 검토 기록, 이벤트 및 메트릭이 포함되어 있습니다. agent_policy에는 policy_lookup에서 사용하는 읽기 전용 동기화 정책이 포함되어 있습니다. 각 Activity는 Workflow에서 사용하는 것과 동일한 결정론적 식별자를 키로 사용하여 레코드를 기록하므로, Lakebase를 Temporal의 재실행 메커니즘의 일부로 만들지 않고도 재시도 후 프로젝션이 따라잡을 수 있습니다.

Unity Catalog는 심사 임계값의 소스입니다. 지속적으로 동기화되는 테이블을 통해 실행 중인 에이전트가 이를 사용할 수 있습니다. 적용된 임계값, 증거 및 후속 사람의 결정은 agent_ops에 기록됩니다. Change Data Feed는 감사 및 분석을 위해 이러한 변경 사항을 Unity Catalog에서 관리하는 기록 테이블에 게시할 수 있습니다.

워커 실패 후 완료된 작업 복구

Temporal은 다른 워커에서 Workflow 상태를 재구축하는 데 필요한 순서가 지정된 Event History를 유지합니다. 이 기록에는 Activity 예약 및 결과, 타이머, 그리고 Signals가 포함됩니다. Replay는 기록된 Event에 대해 Workflow 코드를 실행하고 현재 턴, 수락된 검토 결정, 토큰 사용량 및 수집된 증거와 같은 변수를 재구성합니다.

기록된 Activity 결과는 Activity를 다시 실행하는 대신 재실행 중에 반환됩니다. 완료된 신용 조회는 완료된 상태로 유지되고, 기록된 모델 응답은 해당 실행에 대한 응답으로 유지됩니다. 워커가 실패했을 때 Activity가 진행 중이었고 Temporal이 완료를 기록하지 못한 경우, Temporal은 다른 시도를 예약할 수 있습니다. 에이전트의 경우, 이는 이미 Event History에 기록된 모델 응답을 보존합니다. 완료가 기록되지 않은 모델 호출은 공급자가 처리를 완료했더라도 다시 실행될 수 있습니다.

재시도 정책은 개별 작업 단위로 할당되며 코드에서 재사용할 수 있습니다. 예시에서 모델 호출 액티비티는 3분의 schedule-to-close 타임아웃 내에서 최대 4번의 시도를 허용합니다. 도구 호출 액티비티는 최대 3번의 시도를 허용하며 60초의 start-to-close 타임아웃을 가집니다. Lakebase 액티비티는 15초의 start-to-close 타임아웃과 함께 최대 5번의 시도를 허용합니다.

외부 영향을 안전하게 반복 가능하도록 만들기

한 가지 위험은 Worker가 액티비티 완료를 보고하기 전에 Lakebase 도구 결과 쓰기가 커밋될 수 있다는 점입니다. 그 사이에 연결이 끊어지면 Temporal에 기록된 결과가 없으므로 다른 시도를 예약합니다. 두 시도 모두 동일한 논리적 쓰기를 나타냅니다.

각 Lakebase 레코드는 안정적인 식별자(stable identity)를 가집니다. run_id은(는) 운영 스키마의 기준이 됩니다. message_id은(는) 메시지를, tool_call_id은(는) 도구 호출을, event_id은(는) 마일스톤을, review_id은(는) 검토 라운드를, decision_id은(는) 검토자 명령을 식별합니다. Postgres 기본 키와 고유 제약 조건이 이러한 식별자를 강제합니다.

도구 시작 쓰기는 안정적인 식별자와 최종 상태 가드(terminal-state guard)를 모두 보여줍니다.

재시도는 동일한 tool_call_id을(를) 대상으로 합니다. 최종 조건자(predicate)는 기존의 최종 상태가 아닌 행만 started 상태로 다시 쓰여지도록 허용합니다. 행이 이미 성공 또는 실패 상태인 경우, PostgreSQL은 0개의 행에 영향을 미칩니다. 오류를 발생시키지는 않습니다.

호출자는 0개 행 결과를 검사해야 합니다. LakebaseWriteResult은(는) 영향을 받은 행 수를 반환하지만, 현재 액티비티 래퍼는 0을 실패로 처리하지 않습니다. 프로덕션 코드는 저장된 최종 상태를 확인한 후에만 0을 예상된 무작업(no-op)으로 분류해야 합니다. 그렇지 않으면 충돌을 발생시키거나 기록해야 합니다. 가드된 실행 및 검토 전환에도 동일한 규칙이 적용됩니다.

유사한 upsert가 메시지, 도구 결과 및 Event에 적용됩니다. 결정론적 ID 덕분에 재시도가 동일한 논리적 행으로 수렴되는 한편, 각 가드된 쓰기는 어떤 상태 전환이 유효한지 정의합니다. API는 쓰기가 재시도되는 동안 일시적으로 이전 상태를 표시할 수 있습니다. 액티비티가 성공하면 수락된 행을 조회할 수 있습니다.

부작용을 유발하는 모든 도구에는 이에 상응하는 계약이 필요합니다. 결제 API는 멱등성 키를, 이메일 서비스는 호출자가 제공한 메시지 ID를, 데이터베이스는 고유 제약 조건을 허용할 수 있습니다. 외부 시스템이 중복 제거 메커니즘을 제공하지 않는 경우, 액티비티에는 자체 기록 또는 조정 프로세스가 필요합니다. Temporal은 재시도 시점을 결정합니다. 액티비티는 외부 시스템이 해당 재시도를 처리하는 방법을 결정합니다.

현재 상태 및 운영 메트릭 노출

Event History는 실행 의미론(execution semantics)과 디버깅 세부 정보를 제공합니다. 애플리케이션에는 현재 실행에 대한 인덱싱된 관계형 쿼리가 필요합니다. 사용자 및 상태별 케이스 목록 표시, 증거가 포함된 트랜스크립트 하나 로드, 특정 사용자를 대기 중인 검토 찾기, 여러 실행에 걸친 측정값 집계 등이 이에 해당합니다.

Lakebase는 해당 애플리케이션 뷰를 정규화된 Postgres 스키마에 저장합니다. agent_runs은(는) 현재 상태, Workflow ID, 요청, 총 토큰 수, 타임스탬프 및 추천 메타데이터를 보유합니다. agent_messages은(는) 트랜스크립트를 보존합니다. agent_tool_calls은(는) 인수, 상태, 구조화된 결과, 오류 및 타이밍을 기록합니다. agent_review_decisions은(는) 추천을 안정적인 검토 ID, 검토자 명령, 근거 및 결정 시간과 연결합니다.

또한 이 스키마는 Workflow, 턴(turn), 액티비티 시도 수준에서 명명된 Event 및 메트릭을 기록합니다. FastAPI는 이러한 테이블을 기반으로 하는 run-detail, workflow-metrics, retry-metrics 엔드포인트를 노출합니다. UI는 증거를 수집 중인 실행, 검토를 대기 중인 실행, 실패한 도구를 재시도 중인 실행을 각각 보여줄 수 있습니다. 운영자는 SQL을 사용하여 동일한 행을 쿼리할 수 있습니다.

증거는 Workflow가 완료되기 전에 사용할 수 있습니다. policy_lookup이(가) 완료되면 구조화된 결과가 도구 호출과 함께 저장됩니다. 실행이 AWAITING_REVIEW에 도달하면 언더라이터는 신용 점수, DTI, 임계값, 규칙 결과, 근거 및 추천을 생성한 정책 소스를 볼 수 있습니다.

수동 검토의 내구성 유지 및 오래된 명령 거부

모델이 추천을 반환하면 Workflow는 run_id 및 현재 턴에서 review_id을(를) 파생시킵니다. 대기 중인 검토를 Lakebase에 쓰고, agent.review_pending 이벤트를 기록하고, 프로젝션을 AWAITING_REVIEW로 설정한 다음, workflow.wait_condition을(를) 호출합니다. Temporal은 Worker 프로세스를 점유하지 않고 열려 있는 Workflow를 유지합니다.

API는 언더라이터의 작업을 Signal로 보냅니다. 이를 보내기 전에 API는 Lakebase에 실행이 검토 대기 중으로 표시되는지, 제출된 review_id가 현재 라운드와 일치하는지 확인합니다. 두 확인 중 하나라도 실패하면 API는 충돌을 반환합니다. Workflow는 자체 상태에 대해 명령을 독립적으로 검증하고 오래되거나 중복된 결정을 무시하므로, Lakebase 프로젝션이 지연되는 경우에도 실행을 보호합니다.

Signal을 수락한 후 Workflow는 멱등성 있는 Lakebase 액티비티를 통해 결정을 유지합니다. 승인 또는 거부는 실행을 완료합니다. 추가 정보 요청은 프로젝션을 다시 RUNNING으로 변경하고, 검토자의 근거를 사용자 메시지로 추가하며, 다음 턴을 시작합니다. 턴이 변경되었기 때문에 다음 추천은 새로운 review_id을(를) 받게 됩니다.

API의 202 응답은 Temporal이 Signal을 수신했음을 확인합니다. 비즈니스 수락은 Workflow에서 비동기적으로 발생하므로, 명령이 API 사전 검사를 통과하더라도 검토 상태가 변경된 경우 무시될 수 있습니다. 클라이언트는 결과 상태를 관찰하기 위해 Lakebase 프로젝션을 새로 고칩니다.

Worker 재배포 없이 거버넌스 정책 제공

언더라이팅 임계값은 Worker 코드와 독립적으로 변경됩니다. Unity Catalog의 소스 테이블에는 최소 신용 점수, 자동 승인 DTI, 강제 거절 임계값, 정책 이름과 같은 목적별 값이 포함되어 있습니다.

설정 스크립트는 agent_policy.underwriting_policy_limits. policy_lookup이라는 이름의 지속적인 Lakebase 동기화 테이블을 생성하며, 정규화된 대출 목적에 따라 이 읽기 전용 Postgres 사본을 쿼리합니다. 정책 소유자가 Unity Catalog 소스를 업데이트하면 동기화 파이프라인이 변경 사항을 전파하고, 이후 실행 시 Worker나 API 배포 없이 이를 읽습니다.

정책 결과에는 적용된 임계값, 모든 규칙의 실제 값 및 통과/실패 결과, 그리고 소스가 포함됩니다. 데모는 Lakebase가 비활성화되거나 행을 사용할 수 없을 때 픽스처(fixture) 정책으로 폴백할 수 있으며, 해당 경로를 fixture_fallback로 기록합니다. 규제 대상 Workflow는 대신 실패 시 차단(fail closed)될 수 있습니다. 애플리케이션은 이러한 폴백 결정을 명시적으로 내려야 합니다.

Unity Catalog로 운영 변경 사항 반환

리포지토리는 REPLICA IDENTITY FULL을(를) 설정하여 Lakebase Change Data Feed를 위한 각 agent_ops 테이블을 준비합니다. 관리자는 여전히 스키마에 대해 이 기능을 활성화해야 합니다. 그런 다음 Lakebase는 Postgres 미리 쓰기 로그(write-ahead log)에서 삽입, 업데이트, 삭제를 캡처하고 이를 lb_<table>_history 패턴으로 이름이 지정된 Unity Catalog 관리 Delta 기록 테이블에 배치로 기록합니다.

Change Data Feed는 현재 Public Preview 상태이며 약 15초마다 변경 사항을 플러시합니다. 이 간격은 감사 및 분석에 적합하며, UI는 현재 운영 상태를 위해 Lakebase를 직접 쿼리합니다.

기록 테이블은 실행의 정책 소스, 도구 증거, 액티비티 시도, 검토 대기, 추천 및 수동 결정을 재구성할 수 있습니다. 리포지토리는 이 경로에 대해 소스 스키마를 구성하지만 관찰된 엔드투엔드(end-to-end) Change Data Feed 실행은 포함하지 않습니다. 피드를 활성화하고 대상 테이블을 확인하는 작업은 배포 단계로 남아 있습니다.

시스템 운영

배포는 React/FastAPI를 Temporal Worker와 분리합니다. API 복제본(replica)은 요청 부하에 따라 확장되고, Worker는 Workflow 및 액티비티 Task 백로그와 구성된 동시성에 따라 확장됩니다. Lakebase 오토스케일링은 프로젝트 범위 내에서 데이터베이스 컴퓨팅을 조정합니다.

Temporal Cloud 요금은 Action과 활성 및 보존된 Event History 스토리지를 기준으로 하므로, 재시도 빈도와 장기간 열려 있는 기록도 비용에 영향을 미칩니다. 팀은 여전히 워크로드에 맞게 Kubernetes 복제본, Task Queue, 연결 풀(connection-pool) 제한 및 데이터베이스 범위를 설정해야 합니다. 또는 최신 오픈 소스 릴리스를 사용하여 자체 오픈 소스 Temporal Service를 구축할 수도 있습니다.

Lakebase 클라이언트는 OAuth M2M(machine-to-machine) 인증을 사용합니다. Databricks OAuth 토큰 및 생성된 데이터베이스 자격 증명은 만료되므로, 클라이언트는 1시간의 데이터베이스 자격 증명이 만료되기 전에 SQLAlchemy 연결 풀을 새로 고칩니다. 연결은 TLS를 사용합니다. 순환(rotation)이 없으면 장시간 실행되는 Worker는 예측 가능한 일정에 따라 데이터베이스 실패를 겪게 됩니다.

운영자는 Temporal을 사용하여 Workflow 및 액티비티 기록을 검사하고, Lakebase를 사용하여 애플리케이션 상태 및 메트릭을 쿼리하며, Kubernetes를 사용하여 프로세스 및 배포 상태를 확인합니다. 이를 통해 운영자는 의도적인 검토 대기와 액티비티 재시도, 데이터베이스 액세스 실패 또는 실패한 도구를 구분할 수 있습니다.

증거 및 제한 사항

테스트 스위트에는 워크플로 시퀀싱, 검토 동작, OAuth 연결 구성, 멱등성 영속성, 메트릭 계약, API 워크플로 시작 및 워커 설정에 대해 통과한 21개의 테스트가 포함되어 있습니다. 크래시 복구 스크립트는 결정론적 프로바이더를 사용한 프로세스 실패 실습을 추가합니다.

신청자 및 프로바이더 데이터는 픽스처입니다. 이 리포지토리는 대출 모델, 규제 준수, 프로덕션 보안 제어, 지역별 가용성 또는 대규모 성능을 검증하지 않습니다. 로컬 크래시 실습은 Lakebase가 비활성화된 상태에서 실행되었으므로 Temporal 복구를 격리합니다. 대상 Databricks 환경에서 Change Data Feed를 활성화하고 검증해야 합니다.

다음 단계

더 자세히 알아보고 싶으신가요? 데모를 직접 실행하여 지속성 있는 실행이 작동하는 모습을 확인해 보세요. 참조 구현을 실행하고, 실행 중간에 워커를 중지한 후 에이전트가 복구되는 과정을 지켜보세요. Lakebase를 연결하여 각 결정의 기반이 되는 증거, 정책 및 휴먼 리뷰 상태를 탐색해 보세요.

(이 글은 AI의 도움을 받아 번역되었습니다. 원문이 궁금하시다면 여기를 클릭해 주세요)

최신 게시물을 이메일로 받아보세요

블로그를 구독하고 최신 게시물을 이메일로 받아보세요.