AI 에이전트에 파이프라인 계보 메모리 부여: Cognee를 위한 OpenLineage 커넥터를 구축한 방법
요약
본 글은 자율 AI 에이전트가 데이터의 운영 계보 메모리(operational lineage memory) 부족으로 겪는 한계를 지적합니다. 이를 해결하기 위해, OpenLineage 및 Marquez 표준을 활용하여 Cognee에 커넥터를 구축하는 방법을 상세히 설명하며, 엔지니어링 설계와 인지 기억 추출 메커니즘을 제시합니다.
핵심 포인트
- AI 에이전트는 데이터의 운영 계보(operational lineage)를 알지 못해 한계가 있습니다.
- OpenLineage 표준은 관측 데이터 계보를 추적하는 산업 표준입니다.
- Cognee에 OpenLineage 커넥터를 구축하여 AI 메모리 그래프로 연결했습니다.
- 엔지니어링 설계와 스냅샷 조정 패턴을 통해 구현 방법을 안내합니다.
AI 에이전트에 파이프라인 계보 메모리 부여: Cognee를 위한 OpenLineage 커넥터를 구축한 방법
저자: Soumyajit Ghosh (@somuai)
대상 출판물: Dev.to / Substack
해커톤: Mergetober (WeMakeDevs x Cognee)
풀 리퀘스트: topoteretes/cognee-community#351 (Closes topoteretes/cognee#5554)
1. 새벽 2시의 사건: LLM이 파이프라인 계보 메모리를 필요로 하는 이유
다음 시나리오를 상상해 보세요. 시간은 새벽 2시 15분이고, 데이터 플랫폼의 자동 이상 감지기(anomaly detector)가 P0 페이지를 울립니다. 핵심 임원 매출 대시보드가 비어 있습니다.
내부 AI 운영 어시스턴트를 실행하고 다음과 같이 질문합니다:
"
daily_revenue_summary테이블이 비어 있는데, 상류에서 무엇이 실패했나요?"
만약 당신의 어시스턴트가 전통적인 RAG(Retrieval-Augmented Generation) 또는 단순한 관계형 데이터베이스 커넥터로 백업되어 있다면, 즉각적인 막다른 골목에 부딪힐 것입니다. 에이전트는 daily_revenue_summary의 스키마는 알 수 있지만, 그것을 생성한 파이프라인에 대한 운영 가시성(operational visibility)은 전혀 없습니다. 다음 사실들을 볼 수 없다는 것입니다:
- 상류 Apache Spark 작업(
etl_clean_payments)이 45분 전에OutOfMemoryError를 발생시켰다. - 그 Spark 작업은 외부 마이크로서비스에 의해 스키마가 변경된 Apache Kafka 토픽에서 데이터를 읽는다.
- 하류 dbt 모델이 의존성(dependency)이 완료되어 쓰기 작업을 수행하지 못했기 때문에 건너뛰어졌다.
대부분의 자율 AI 에이전트는 근본적인 아키텍처적 사각지대에 시달립니다: 그들은 운영 계보 메모리(operational lineage memory)를 전혀 가지고 있지 않습니다. 그들은 데이터가 저장된 상태에서 어떤 모습인지 알지만, 데이터가 분산 컴퓨팅 엔진을 거쳐 어떻게 이동하고, 변환되거나, 실패하는지에 대한 인지 지도(cognitive map)는 없습니다.
이 문제를 해결하기 위해, 우리는 OpenLineage 및 Marquez 데이터 소스 커넥터를 Cognee를 위해 설계하고 구축했습니다(topoteretes/cognee-community#351), 오픈 메타데이터 표준을 자율 AI 메모리 그래프로 연결합니다.
이 심층 분석에서는 엔지니어링 설계, 인지 기억 추출 메커니즘, 전체 스냅샷 조정 패턴(full-snapshot reconciliation pattern), 그리고 이를 프로덕션 환경에서 실행하는 방법을 안내합니다.
2. 기반 이해하기: OpenLineage & Cognee
코드에 들어가기 전에, 이 통합을 구동하는 두 시스템을 살펴보겠습니다.
OpenLineage 표준 및 Marquez
Linux Foundation (LF AI & Data) 산하에서 생성된 OpenLineage는 관측 데이터 계보(observational data lineage)를 위한 오픈 소스 산업 표준입니다. OpenLineage는 다음 항목들을 추적하기 위한 확장 가능한 명세(specification)를 정의합니다:
- Jobs: 계산 단위 (Airflow 태스크, Spark 애플리케이션, dbt 모델, Flink 잡).
- Datasets: 소비되거나 생성된 데이터 저장소 (PostgreSQL 테이블, Iceberg 데이터셋, S3 Parquet 경로).
- Runs: 상태 전환(
START,RUNNING,COMPLETE,FAIL,ABORT)을 가진 작업의 특정 실행. - Facets: 스키마 정의, SQL 쿼리 텍스트, 데이터 품질 단언(data quality assertions), 실패 스택 트레이스 등을 포함하는 세밀한 메타데이터 첨부물.
Marquez는 OpenLineage를 위한 참조 백엔드 구현체 역할을 하며, 깔끔한 REST API를 통해 계보 그래프를 저장하고 노출합니다.
Cognee란 무엇인가?
Cognee (topoteretes/cognee)는 AI 에이전트를 위해 특별히 설계된 오픈 소스 메모리 엔진입니다. Cognee는 정보를 단절된 벡터 청크로 취급하는 대신 다음 작업을 수행합니다:
dlt(data load tool) 기반의 선언적 데이터 로딩 파이프라인을 통해 멀티모달 데이터를 **수집(Ingests)**합니다.- 들어오는 문서를 **인지화(Cognifies)**합니다: 언어 모델을 사용하여 개념, 유형화된 엔티티, 그리고 엔티티 간 관계를 추출합니다.
- 결과로 나온 지식 그래프(FalkorDB, Neo4j 또는 NetworkX 기반)와 벡터 인덱스 검색(LanceDB, Qdrant)을 결합하여 **영구 저장(Persists)**합니다.
- 그래프 완성 검색 실행: LLM 에이전트가 벡터 거리만으로는 해결할 수 없는 복잡한 질문에 답하기 위해 깊은 그래프 경로를 탐색하도록 합니다.
3. 아키텍처 및 데이터 흐름: 파이프라인 텔레메트리에서 인지 그래프까지
원시 텔레메트리를 LLM 메모리 그래프에 연결하려면 체계적인 수집 파이프라인이 필요합니다. 아래 다이어그램은 운영 오케스트레이터에서 Cognee로 계보(lineage) 메타데이터가 어떻게 흐르는지 보여줍니다:

아키텍처적 문제: 원시 JSON 수집 vs. 의미론적 투영(Semantic Projections)
기업 환경의 OpenLineage 배포는 매일 수백만 건의 원시 JSON RunEvent 페이로드를 수집합니다. 일반적인 원시 이벤트 페이로드는 소켓 주소, 프로듀서 클라이언트 버전, 하트비트 카운터 및 시스템 메트릭 등 500줄 이상의 중첩된 전송(transport) 텔레메트리를 포함합니다.
원시 JSON RunEvent를 LLM 컨텍스트에 그대로 덤프하려고 하면 심각한 문제가 발생합니다:
- 컨텍스트 창 과부하 (Context Window Flooding): 단일 파이프라인 실행만으로도 실질적인 통찰력을 제공하지 못하면서 수만 개의 토큰을 소모할 수 있습니다.
- 엔티티 추출 성능 저하 (Degraded Entity Extraction): 일반 목적 LLM은 깊게 중첩된 JSON 보일러플레이트 속에 묻힌 그래프 관계를 추론하는 데 어려움을 겪습니다.
- 보안 취약점 (Security Vulnerabilities): 원시 이벤트 페이츠에는 종종 마스킹되지 않은 데이터베이스 연결 URI, 서비스 계정 토큰 또는 환경 매개변수가 포함되어 있습니다.
우리의 해결책: 고밀도 의미론적 Markdown 투영(High-Density Semantic Markdown Projections)
원시 JSON 이벤트를 스트리밍하는 대신, 저희 커넥터는 Marquez 카탈로그(/api/v1/namespaces, /api/v1/jobs, /api/v1/runs, /api/v1/datasets)를 쿼리하고 각 작업을 위한 고밀도의 구조화된 Markdown 문서를 합성합니다:
# OpenLineage Job: analytics.etl_clean_payments
- **Namespace**: analytics
...
Cognee가 이 투영(projection)을 cognify 엔진을 통해 처리할 때, LLM은 즉시 관계를 인식합니다:
- 엔티티
etl_clean_payments가raw_transactions로부터 읽어옴 (READS_FROM). - 엔티티
etl_clean_payments가daily_revenue_summary로 작성함 (WRITES_TO). - 엔티티
etl_clean_payments는 상태가FAIL이고 오류가OutOfMemoryError인run_84920을 통해 실행됨 (HAS_RUN).
4. 커넥터 구현: 문서 모드 및 방어적 엔지니어링
packages/connector/openlineage/cognee_community_connector_openlineage/openlineage.py의 구현 세부 사항을 살펴보겠습니다.

A. Cognee 문서 모드를 통한 라우팅
Cognee에서는 구조화된 테이블 데이터가 일반적으로 관계형 SQL 스키마로 라우팅됩니다. 하지만 운영 계보 그래프 모델은 비정형 문서 파이프라인을 통해 수집될 때 가장 풍부합니다.
저희는 dlt 소스에 Cognee의 내부 마커를 태그합니다:
from cognee.tasks.ingestion.dlt_utils import DOCUMENT_SOURCE_ATTR
source = _openlineage(
...
이는 Cognee가 생성된 계보 문서를 엔티티 추출로 직접 라우팅하여 작업(jobs), 데이터셋(datasets), 스키마(schemas), 실행 기록(run histories)에 대한 타입이 지정된 노드를 구성하도록 보장합니다.
B. 유령 지식 방지: 전체 스냅샷 교체 (`write_disposition=
- 모든 동기(sync) 실행은 계보 카탈로그의 권위 있는 상태를 검색합니다.
write_disposition="replace"는 스테이징 테이블을 현재 상태로 대체합니다.- Marquez에서 오래된 작업(job)이 삭제되면, 스테이징에서도 사라집니다.
- Cognee의 내부
orphan_cleanup은 사라진 엔티티를 감지하고 지식 그래프와 벡터 인덱스 모두에서 이를 정리(reconcile)합니다. - 변경되지 않은 작업은 안정적인 콘텐츠 해시를 유지하여, 불필요한 재-인지화(re-cognification)를 방지하고 LLM API 토큰을 절약합니다.
C. 방어적 비밀 값 제거 (Defensive Secret Sanitization)
계보 속성(Lineage facets)은 데이터베이스 자격 증명, JDBC 비밀번호 또는 연결 문자열이나 실행 매개변수에 포함된 Bearer 토큰을 실수로 노출할 수 있습니다.
이 커넥터는 문서 발행 전에 방어적 정규 표현식 마스킹을 적용합니다:
_SECRET_PATTERN = re.compile(
r'(?i)(token|secret|password|passwd|api[_-]?key|access[_-]?key|auth|bearer)\s*[:=]\s*["\']?([^"\\'\s]+)["\']?'
)
...
어떠한 원시 자격 증명도 지식 그래프나 LLM의 프롬프트에 도달하지 않습니다.
5. 실시간 실행 및 테스트 검증 (Live Execution & Test Verification)
모든 커넥터 코드는 네임스페이스 검색, 속성 파싱, 비밀 값 제거(secret redaction), 스냅샷 대체 등을 포괄하는 엔드투엔드 테스트로 검증되었습니다:

테스트 스위트 하이라이트 (tests/test_openlineage.py):
- 10개 중 10개 단위 및 통합 테스트 통과: 오프라인에서 7.62초 동안 실행되었습니다.
- 제로 네트워크 불안정성(Zero Network Flakiness):
FakeMarquezClient를 구현하여, 외부 Docker 서비스나 클라우드 종속성을 구동하지 않고도 GitHub Actions의 CI 러너가 페이지네이션, 속성 파싱 및 오류 백오프 테스트를 수행할 수 있게 합니다. - Ruff 준수:
ruff check와ruff format모두에서 100% 깨끗합니다.
6. 실제 사례 분석: 그래프 완성으로 실패 진단하기 (Real-World Walkthrough: Diagnosing Failures with Graph Completion)
AI 에이전트가 이 커넥터를 프로덕션 환경에서 어떻게 사용하는지 살펴보겠습니다.
1단계: 계보 토폴로지 수집 (Ingesting Lineage Topologies)
import asyncio
import os
import cognee
...
2단계: 자율 에이전트를 이용한 그래프 질의 (Querying the Graph with an Autonomous Agent)
메모리 그래프가 채워지면, 에이전트는 복잡한 운영 질문에 답할 수 있습니다:
async def query_agent():
question = (
"daily_revenue_summary에 쓰는 상위 작업(upstream jobs)은 무엇인가?", "...
에이전트 응답 및 계보 탐색 (Agent Response & Lineage Traversal)
다음은 정확한 질의 입력과 합성된 그래프 응답입니다:

에이전트가 그래프를 탐색한 방식을 주목하세요:
- 노드
daily_revenue_summary를 찾았습니다. - 들어오는
WRITES_TO엣지를 역추적하여etl_clean_payments를 발견했습니다. HAS_RUN관계를 검사하여 실행run_84920을 표시했습니다.- 정확한 오류 세부 정보(
OutOfMemoryError: Java heap space)를 가져와 엔지니어에게 대시보드에 미치는 하위 영향(downstream impact)에 대해 경고했습니다.
전체 추론 과정은 2초도 안 되는 시간에 순수하게 구조화된 지식 그래프에서 파생되었습니다.
7. 주요 엔지니어링 시사점 (Key Engineering Takeaways)
Mergetober를 위해 이 통합을 구축하면서, 에이전트 기반 데이터 인프라(agentic data infrastructure)에 대한 세 가지 핵심 아키텍처 교훈을 얻었습니다:
AI 자동 생성 콘텐츠
본 콘텐츠는 Dev.to AI tag의 원문을 AI가 자동으로 요약·번역·분석한 것입니다. 원 저작권은 원저작자에게 있으며, 정확한 내용은 반드시 원문을 확인해 주세요.
원문 바로가기