Загружаем каталог…
Загружаем каталог…
이 시리즈는 Databricks 환경에서 구축했던 항공/운송/물류 분야의 RAG Agent 구축 프로젝트 관련 내용을 정리하는 시리즈입니다. 요구사항과 관련하여 Databricks의 어떤 기술을 사용해 이를 충족하였으며, 수행 중 발생한 문제 상황을 어떻게 해결했는지 정리합니다. 실제 고객사에서 수행한 프로젝트와는 별도의 내용으로, 기술적인 내용만을 정리합니다. 1. 데이터 수집 개요( Data Ingestion ) 데이터 수집(Data Ingestion) 작업은 외부 데이터 소스로부터 최신 문서를 주기적으로 스캔하여, 변경되거나 신규로 추가된 문서를 감지하고 Databricks 환경(Unity Catalog 및 Managed Volume)에 안전하게 파일 형태로 다운로드 및 적재하는 파이프라인의 첫 번째 단계입니다. 이 파이프라인은 Databricks Workflow 상에서 각각 독립적인 Task 노트북으로 구성되어 동작합니다. 1.1. 전체 데이터 수집 및 적재 흐름도 1.2. 주요 태스크(Task) 정의 태스크명 역할 및 핵심 기능 Task #0. Initialize 파이프라인 구동에 필수적인 Delta Table(메타데이터, 이벤트 로그 등)을 부트스트랩하고 최신 스키마 상태로 유지 관리합니다. Task #1. Metadata Sync Source Storage를 스캔하여 현재 파일 목록의 스냅샷을 생성하고, 변경 감지 대상이 되는 Delta Table에 병합하여 변경 이력(CDF)을 발생시킵니다. Task #2. Content Export Task #1에서 발생시킨 변경 범위(Delta Version)를 기준으로 실제 본문 콘텐츠와 이미지 데이터를 추출하여 Managed Volume 경로에 병렬 적재합니다. 2. [Task #0] 파이프라인 환경 초기화 및 리셋(Pipeline Initialize) 수집 프로세스가 시작되기 전 필요한 스토리지 스키마를 동기화하고, 데이터 분석 오류나 수집 장애 등으로 인해 전체 인덱싱이 필요할 때 안전하고 영속적인 리셋을 수행하는 메커니즘을 다룹니다. Delta Lake와 Delta Table Databricks의 표준 고성능 테이블 포맷입니다. 일반적인 관계형 데이터베이스(RDB)나 Parquet 파일과 달리 트랜잭션 로그를 생성하여 데이터 무결성을 보장하고, 테이블 버전별 타임 트래블(이전 상태 복구) 및 변경 사항만 정밀 추적하는 기능(Change Data Feed)을 지원합니다. 2.1. bootstrap.py 모듈 및 리소스 부트스트랩 데이터 적재를 위한 Delta 테이블들의 최초 생성 및 구조 갱신에 관련된 책임을 지는 모듈입니다. 2.1.1. ensure_all(spark, names) 역할 : 파이프라인 실행에 필수적인 모든 Delta 테이블과 뷰를 자동으로 생성하는 메인 오케스트레이터입니다. 동작 로직 테이블 정보가 정의된 Names 객체( names )를 참조하여 내부 함수인 ensure_file_meta() , ensure_events() , ensure_log_run() 등을 순차적으로 호출합니다. 테이블 생성이 완료된 후에는 _ensure_columns() 함수를 이용해 최신 소스코드 스키마 정의와 데이터베이스 내 실물 스키마 간의 싱크를 강제로 유지합니다. 2.1.2. ensure_file_meta(spark, names) 역할 : 원본 파일의 메타데이터와 최종 상태 스냅샷을 보관하는 raws.file_meta 테이블을 생성합니다. 동작 로직 DeltaTable.createIfNotExists() API를 사용해 물리적인 테이블을 생성합니다. 테이블 속성(Properties)으로 delta.enableChangeDataFeed = true 를 주입하여 생성합니다. Change Data Feed (CDF) 기능이란? 데이터베이스 테이블에 발생하는 모든 변화(추가, 수정, 삭제)를 마치 Git 커밋 로그처럼 차례대로 기록해 두는 기능입니다. 이를 활성화해 두어야만 Task #2 단계에서 "지난 수집 주기 이후 새로 추가되거나 바뀐 문서만" 정확히 구분하여 효율적으로 다운로드할 수 있습니다. 핵심 관리 테이블 구조 및 스키마 요약 파이프라인 구동 및 관리를 위해 사용되는 3가지 핵심 Delta 테이블의 구조는 다음과 같습니다. 테이블 논리명 물리 테이블 경로 데이터 보존/적재 방식 핵심 속성 및 파티션 원천 파일 메타데이터 raws.file_meta 현재 최신 스냅샷 상태 유지 (MERGE) CDF 활성화 ( enableChangeDataFeed=true ) 수집 완료 이벤트 _meta.file_export_events 수집 성공/삭제 이벤트 누적 (Append-only) 파티션 키: event_date 수집 주기 실행 이력 _meta.cdf_export_run 수집 주기(Cycle)별 런타임 결과 적재 실행 메트릭 및 에러 추적 2.1.3. ensure_events(spark, names) / ensure_log_run(spark, names) 역할 : 작업 성공 로그 적재용 _meta.file_export_events 및 실행 상태 로깅용 _meta.cdf_export_run 테이블을 생성합니다. 동작 로직 ensure_events 는 파티션 키로 event_date 컬럼을 주입하여 생성함으로써 날짜별 쿼리 성능을 최적화합니다. 데이터 형식뿐만 아니라 스키마 정의에 기술된 주석(Comment) 정보까지 보존하여 데이터 카탈로그 명세서를 실시간으로 동기화합니다. 2.1.4. _ensure_columns(spark, table_name, schema) 역할 : 소스코드에 새로운 nullable 필드가 추가되는 등 설계 사양이 변경되었을 때, 기존의 테이블에 Alter 조작을 수행하는 스키마 진화(Schema Evolution) 도구입니다. 동작 로직 Spark API를 통해 대상 테이블의 실물 필드 세트를 가져와 schemas.py 에 선언된 스키마 필드 명세와 대조(Difference)합니다. 신규 컬럼이 발견되면 누락된 필드의 타입(DataType)과 코멘트 정보를 조합하여 ALTER TABLE {table_name} ADD COLUMNS (...) SQL 쿼리를 동적으로 조립하고 스파크 세션을 통해 실행합니다. # bootstrap.py - 스키마 동적 확장을 위한 ALTER SQL 실행 로직 col_defs = [] for f in missing: comment = (f.metadata.get("comment") or "").replace("'", "''") col_defs.append(f"`{f.name}` {f.dataType.simpleString()} COMMENT '{comment}'") spark.sql(f"ALTER TABLE {table_name} ADD COLUMNS ({', '.join(col_defs)})") 2.2. 전체 수집 및 인덱싱 리셋( full_reset.py ) TRUNCATE와 DELETE의 역할 분할 TRUNCATE (구조 유지 내용 소거) : 파싱( _parsed ), 청킹( _chunked ), 벡터 데이터( _vectorized ) 테이블의 경우 테이블의 권한 정보, 주석 메타데이터, 그리고 연동되어 작동 중인 AI Search Index 서비스가 깨지지 않도록 물리 구조는 그대로 두고 알맹이 데이터만 비우는 TRUNCATE TABLE 을 실행합니다. DELETE (선택적 소거) : 메타 테이블( file_meta , file_export_events )은 여러 카테고리의 데이터가 공존하므로, 특정 subject 범위만 골라 DELETE WHERE category IN (...) 구문을 수행해 대상을 정확히 구분하여 소거합니다. 안전장치 관리자의 실수로 인한 데이터 유실을 방지하기 위해 변수 파라미터 confirm 의 기본값을 "false" 로 고정합니다. 노트북 실행부 첫 머리에서 해당 파라미터가 정확히 "true" 문자열로 매칭되지 않을 경우 dbutils.notebook.exit() 를 호출하여 실행을 즉각 차단합니다. 3. [Task #1] 소스 스토리지 메타데이터 동기화(Metadata Sync) 지정된 소스 스토리지 내의 폴더 구조와 파일 목록을 조회하여 이전에 수집된 목록과 대조하고, 수정 및 신규 삽입에 해당하는 데이터 변경을 Delta Change Data Feed(CDF) 상에 각인하는 흐름입니다. 3.1. API 인증 및 크리덴셜 관리 ( SourceServiceFactory ) 3.1.1. from_sa_secret(dbutils, secret_scope, secret_key, delegate_email) 역할 : Databricks Secret Manager로부터 암호화된 소스 스토리지 어카운트의 프라이빗 키 JSON을 안전하게 꺼내와, 이를 기반으로 구글 클라이언트 라이브러리에 연동 가능한 인증 개체 및 서비스를 구성합니다. 동작 로직 : dbutils.secrets.get 메서드로 암호 키 페이로드를 가져온 뒤 JSON 역직렬화를 거쳐 Google Credentials 객체를 초기화합니다. 3.2. 수집 대상 디렉토리 스캔( FolderListingSource ) 3.2.1. FolderListingSource.fetch_snapshot(spark) 역할 : 소스 스토리지에 분산되어 있는 대상 디렉토리 ID들을 전체 조회하여 표준 수집 컬럼 규격에 맞는 단일 Spark DataFrame을 빌드합니다. 동작 로직 노트북 레벨에서 주입된 folder_config.json 을 읽어 active 상태인 폴더 ID 목록을 획득합니다. 각 폴더 정보(ID, Category, Subcategory)를 토대로 루프를 돌며 해당 폴더의 직속 파일 객체를 순차적으로 취합합니다. 획득된 파일 목록을 _rows_to_df() 헬퍼 메서드를 통해 file_meta_schema 스키마 규격으로 정규화하여 스파크 데이터프레임화합니다. 이 과정에서 각 파일의 고유 메타 정보에 카테고리 태그 및 한국 표준시(KST)로 변환된 변경 시점 시각 정보( modified_at )를 병합합니다. 3.2.2. StorageClient.list_files_in_folder(folder_id, mime_types) 역할 : 특정 폴더의 직속 하위 파일들을 페이지네이션을 처리하며 스캔합니다. 동작 로직 : 전달받은 mime_types 목록을 OR 구문으로 묶어 드라이브 조회 쿼리 스트링( q )을 작성합니다. ( 'folder_id' in parents and trashed=false and (mimeType='...' or ...) ) supportsAllDrives=True 설정을 통해 공유 드라이브를 포함한 모든 클라우드 저장 공간을 빠짐없이 스캔하며, nextPageToken 을 추적하여 단 하나의 누락 문서 없이 모든 하위 파일 데이터 구조를 리턴합니다. 3.3. 메타데이터 병합 적재( FileMetaMerger ) 3.3.1. merge(other_df) 역할 : 스토리지 스캔 데이터( other_df )를 기존의 메타 테이블( names.file_meta )에 물리적인 MERGE INTO 방식으로 반영하고, 그 이력 버전을 추적합니다. 핵심 병합 조건 및 로직 식별 키 조건 : t.file_id = s.file_id (기존 DB 테이블 t 와 소스 데이터프레임 s 매핑) 수집 격리 범위 조건( _scope_predicate ) : 카탈로그 데이터 정합성을 해치지 않기 위해 t.category IN (subject_categories) 조건을 병합 타깃 범위로 강제 지정합니다. 이를 통해 해당 영역의 데이터만 수정 및 매칭 비교 프로세스가 일어나며, 동일 테이블을 공유하는 다른 subject 행은 안전하게 보호됩니다. 고아 대상 자동 소거( whenNotMatchedBySourceDelete ) : 드라이브에서 삭제되었거나 타 폴더로 옮겨가면서 이번 스냅샷 소스 목록에서 제외된 수집 범위 내의 파일이 감지될 경우, Delta Engine에 의해 자동으로 감지되어 메타 데이터베이스에서 물리 삭제( DELETE )됩니다. # file_meta_merger.py - Delta Lake 병합(MERGE INTO) 조작부 (file_meta_dt.alias("t") .merge(source=other_df.alias("s"), condition=f"t.file_id = s.file_id AND {scope}") .whenMatchedUpdate( condition="t.modified_at < s.modified_at", set=self._update_set(other_df) ) .whenNotMatchedInsert(values=self._insert_set(other_df)) .whenNotMatchedBySourceDelete(condition=scope) .execute()) 메타데이터 테이블 MERGE 처리 조건표 변경 시나리오 병합 분기 조건 (Merge Clause) 결과 change_type DB 물리 동작 설명 신규 파일 생성 WHEN NOT MATCHED THEN INSERT INSERT 새로운 file_id 를 가진 행을 raws.file_meta 에 추가 삽입 기존 파일 내용 수정 WHEN MATCHED AND (t.modified_at != s.modified_at) THEN UPDATE UPDATE 원본의 수정 시각 정보와 다를 경우, 파일명 및 modified_at 등을 갱신 파일 삭제 또는 폴더 이탈 WHEN NOT MATCHED BY SOURCE AND (t.category IN (...)) THEN DELETE DELETE 수집 스코프 내 에 포함되어 있으나 업로드 스냅샷 소스에서 사라진 대상을 감지해 삭제 3.3.2. _collect_result(dt, pre_version) 역할 : 병합 작업 완료 후 발생한 Delta Transaction History 내역을 역추적하여 수집 통계를 집계하고 병합 결과를 반환합니다. 동작 로직 대상 Delta Table의 메타 히스토리 정보를 필터링하여 이전에 획득했던 pre_version 보다 큰 트랜잭션 버전을 조회합니다. 트랜잭션 내부의 operationMetrics 딕셔너리에서 numTargetRowsInserted , numTargetRowsUpdated , numTargetRowsDeleted 항목을 추출하여 가산(Sum)해 낸 뒤 이를 MergeResult 구조화된 데이터 클래스 객체로 가공해 리턴합니다. 3.4. MergeTaskValues (상태값 전파) set_from_merge_result() 메서드를 호출하여 FileMetaMerger 의 결과값( MergeResult )을 Databricks Task Values에 저장합니다. start_v 와 end_v 버전 번호가 전파되어 후속 다운로드 작업(Task #2)이 대상 테이블 전체를 풀 스캔하는 대신, 지정된 최소 Delta 버전을 타깃으로 하여 증분 변경 내역(CDF)만 선별적으로 다운로드할 수 있게 조율합니다. 4. [Task #2] Source Data 다운로드 및 Volume 적재(Content Export) 이전 단계인 Task #1에서 전파된 버전 정보를 추적하여 실제 변경사항이 존재하는 문서들의 콘텐츠를 직접 내려받고, 병렬 스레드를 사용해 최적화된 경로의 Databricks Volume에 파일로 저장합니다. 4.1. CDF 기반 변경분 추출( FileMetaCdfReader ) 4.1.1. read_changes(start_v, end_v) 역할 : 지정된 Delta Version 범위 내에서 발생한 테이블의 실시간 이벤트 로그를 Change Data Feed에서 가공 처리합니다. 동작 로직 .option("readChangeFeed", "true") 옵션을 활성화하여 Delta 테이블을 로드합니다. 작업 중 파일 상태 갱신 과정에서 중복 수집 및 불필요 로직 수행을 방지하기 위해 변경 타입( _change_type )이 insert , delete , update_postimage 에 해당하는 변화 시점 로그만을 격리 선별한 뒤 데이터프레임을 .cache() 처리합니다. # cdf_reader.py - Change Data Feed 기반 변경 행 조회 return (self._spark.read.format("delta") .option("readChangeFeed", "true") .option("startingVersion", start_v) .option("endingVersion", end_v) .table(tableName=self._table_name) .filter(F.col("_change_type").isin("insert", "delete", "update_postimage")) .cache()) 4.1.2. extract_export_targets(changes) 역할 : CDF 변경 히스토리 중 실제 물리적인 파일 다운로드 조작이 수반되어야 하는 파일 객체만을 별도로 필터링하여 드라이버 메모리에 수집합니다. 동작 로직 삭제( delete ) 이벤트를 배제하고 파일 신규 삽입( insert ) 및 수정 완료본( update_postimage ) 로그에 매칭되는 건들만 추출하여 드라이버 내부의 list[dict] 데이터 형태로 가공 반환합니다. 4.2. 파일 본문 다운로드 및 Volume 저장( Exporter ) 4.2.1. export_to_volume(file_meta, output_dir) 역할 : 특정 문서 파일에 대한 수집 요청 정보를 인계받아 해당 문서의 본문(content)을 다운로드하고 Unity Catalog Volume의 로컬 파일 시스템 레이아웃으로 실체화합니다. 동작 로직 _MIME_HANDLERS 매핑 딕셔너리에 지정된 규칙에 맞춰 문서 타입별 전용 메서드( fetch_document , fetch_spreadsheet , fetch_slide )를 호출하여 본문 전체의 페이로드를 입수합니다. 다운로드 성공 시 Volume 디렉터리 내에 {output_dir}/{folder}/{file_id}/ 형태로 하위 격리 디렉터리를 만들고, 해당 디렉터리 내에 document.json (순수 본문 페이로드)과 document_metadata.json (수집 카테고리 및 URL 등의 파일 메타데이터) 구조로 저장 처리합니다. 이후 ImageExporter 인스턴스를 즉각 구동하여 본문 내 이미지 데이터 저장을 시도합니다. 이때 이미지 저장 처리 도중 예외가 발생하더라도 문서 수집 자체는 영향이 없도록 프로세스를 예외 처리 블록으로 격리합니다. 4.2.2. 재시도 가드 데코레이터( _retry ) 역할 : API 통신 시 흔히 발생하는 일시적인 통신 단절 및 API Quota 초과 현상에 대해 자동 재시도 루프를 수행합니다. 동작 로직 tenacity 라이브러리를 기반으로 구현되어 있습니다. HTTP 상태 코드 중 429 (Too Many Requests), 500 , 502 , 503 , 504 에러나 ConnectionError 와 같은 전송망 에러 발생을 체크( _is_retryable )하여 최대 6회까지 동적 대기 시간 지터(Jittered Exponential Wait) 규칙에 따라 자동으로 지연 및 재실행을 보정 처리합니다. 4.3. 병렬 다운로드 실행 엔진( ParallelExporter ) 4.3.1. run(export_targets) 역할 : 다량의 수집 타깃 목록( export_targets )을 병렬 다중 스레드로 전환하여 다운로드 실행 시간을 최소화합니다. 동작 로직 파이프라인 매개변수로 공급된 max_workers 정밀 값을 가져와 ThreadPoolExecutor 스레드 풀을 선언합니다. 각 태스크 스레드는 드라이브 엑스포터의 export_to_volume 메서드를 참조하여 다운로드를 동시 수행하며, 수집 과정의 성공/실패 여부를 모아 취합한 수집 성공 로그 리스트( file_logs_rows )와 에러 건수를 최종 카운트하여 메인 드라이버에 리턴합니다. # export_runner.py - ThreadPoolExecutor 기반 병렬 다운로드 오케스트레이션 with ThreadPoolExecutor(max_workers=self._max_workers) as executor: futures = [executor.submit(self._exp
То, что RADAR обнаружил и классифицировал для этой возможности. Это опубликованный источником текст, а не подтверждение, что предложение ещё действует.
[DBRX] Databricks RAG Agent - 02. Data Ingestion. 이 시리즈는 Databricks 환경에서 구축했던 항공/운송/물류 분야의 RAG Agent 구축 프로젝트 관련 내용을 정리하는 시리즈입니다. 요구사항과 관련하여 Databricks의 어떤 기술을 사용해 이를 충족하였으며, 수행 중 발생한 문제 상황을 어떻게 해결했는지 정리합니다. 실제 고객사에서 수행한 프로젝트와는 별도의 내용으로, 기술적인 내용만을 정리합니다. 1. 데이터 수집 개요( Data Ingestion ) 데이터 수집(Data Ingestion) 작업은 외부 데이터 소스로부터 최신 문서를 주기적으로 스캔하여, 변경되거나 신규로 추가된 문서를 감지하고 Databricks 환경(Unity Catalog 및…
Открыть источник