์์ง ๐กฒ Lake ๐กฒ Spark ๐กฒ VectorDB ํ๋ฆ ์๋ํ
- 1. ์ค์ผ์คํธ๋ ์ด์ ์ํคํ ์ฒ
- 2. ์ข ํฉ ์ค์ต ์์ ์ฝ๋ 1
- 3. ์ข ํฉ ์ค์ต ์์ ์ฝ๋ 2
1. ์ค์ผ์คํธ๋ ์ด์ ์ํคํ ์ฒ
- โ์์ง ๐กฒ Data Lake ๐กฒ Apache Spark ๐กฒ VectorDBโ๋ก ์ด์ด์ง๋ ํ์ดํ๋ผ์ธ:
- ํ๋์ ์ธ Enterprise AI ๋ฐ ๋๊ท๋ชจ ํ์ด๋ธ๋ฆฌ๋ RAG(Retrieval-Augmented Generation) ์์คํ ์ ํ์ค ๋ฐฑ์๋ ์ํคํ ์ฒ
- ์ค๋ฌด ํ๊ฒฝ์์๋ ํ ๋ผ๋ฐ์ดํธ๊ธ ๋์ฉ๋ ๋น์ ํ ๋ฐ์ดํฐ(๋ฌธ์, ๋ก๊ทธ)๊ฐ ์ ์ ๋๋ฏ๋ก, ๋จ์ผ Airflow ์์ปค์์ ๋ฐ์ดํฐ๋ฅผ ๊ฐ๊ณตํ๋ ๊ฒ์ ๋ถ๊ฐ๋ฅ
- ๋ฐ๋ผ์ Airflow๋ ์ ์ด(Orchestration)๋ง ๋ด๋นํ๊ณ , ์ค์ ์ค๋ ์ฐ์ฐ์ ๋ถ์ฐ ์ปดํจํ ์์ง(Spark)์ ์์ํ๋ ๊ตฌ์กฐ๋ฅผ ์ทจํจ
1.1 ์ ์ ๋ฐ ๊ฐ๋
- ์ค์ผ์คํธ๋ ์ด์
์ํคํ
์ฒ(Orchestration Architecture)
- ๋ฐ์ดํฐ ์์ง๋์ด๋ง ๋ฐ AI ์ธํ๋ผ์์ ๋ถ์ฐ๋ ์์คํ
, ์๋น์ค, ๋ณต์กํ ๋ฐ์ดํฐ ํ์ดํ๋ผ์ธ์
- ์คํ ์์, ์์กด์ฑ, ์์ ํ ๋น, ์์ธ ์ฒ๋ฆฌ ๋ฑ์ ์ค์์์ ์ ๋ฐ์ ์ผ๋ก ์กฐ์จ(Orchestrate)ํ๊ณ ํต์ ํ๋ ํต์ฌ ํ๋ ์์ํฌ๋ฅผ ์๋ฏธ
- ๋จ์ํ โ์ค์ผ์ค๋ฌ์ ๋ง์ถฐ ๋ฐฐ์น ์คํฌ๋ฆฝํธ๋ฅผ ์คํํ๋ ๊ฒโ์ ๋์ด,
- ๊ฑฐ๋ํ ๋ถ์ฐ ์ปดํจํ ํ๊ฒฝ์ ์์ ํ๊ณ ๋ฉฑ๋ฑ์ฑ ์๊ฒ ๊ด๋ฆฌํ๊ธฐ ์ํ ๊ณ ๋์ ์์คํ ์ํคํ ์ฒ
- ์๋ง์ ๋ถ์ฐ ์ธํ๋ผ์ ๋ฐ์ดํฐ ์์ค๋ค ์ฌ์ด์์ ๋ค์์ ๋ํ ํด๋ต์ ์ ์ํ๋ ๊ณ ๋์ ์์คํ
์ค๊ณ ๊ธฐ์
- ์ด๋ป๊ฒ ๊ฒฐํจ์ ๊ฒฉ๋ฆฌํ๊ณ ,
- ์ด๋ป๊ฒ ์์์ ์ต์ ํํ๋ฉฐ,
- ์ด๋ป๊ฒ ๋ฐ์ดํฐ ํ๋ฆ์ ์์ ํ๊ฒ ๋ณด์ฅํ ๊ฒ์ธ๊ฐ
- ๋ฐ์ดํฐ ์์ง๋์ด๋ง ๋ฐ AI ์ธํ๋ผ์์ ๋ถ์ฐ๋ ์์คํ
, ์๋น์ค, ๋ณต์กํ ๋ฐ์ดํฐ ํ์ดํ๋ผ์ธ์
- ์ํคํ
์ฒ ์ ํ: ํจํด ๊ธฐ๋ฐ ๋ถ๋ฅ
-
์ค์ผ์คํธ๋ ์ด์ ์ ๋ณต์กํ ํ๋ถ ์๋น์ค๋ฅผ ์กฐ์จํ๋ ๋ฐฉ์์ ๋ฐ๋ผ ํฌ๊ฒ ๋ ๊ฐ์ง ์ค๊ณ ํจํด์ผ๋ก ๋๋จ
- ์ค์ผ์คํธ๋ ์ด์
ํจํด (Orchestration Pattern) - ์ค์ ์ง์คํ
- ์ค์์ ๊ฐ๋ ฅํ โ์งํ์(Orchestrator)โ ์ญํ ์ ํ๋ ์์ง(์: Apache Airflow, Temporal, Prefect)์ ๋๊ณ ,
-
์ด ์์ง์ด ๋ชจ๋ ์๋น์ค์ ํ์คํฌ์ ์ํ๋ฅผ ์ง์ ์ ์ดํ๊ณ ๋ช ๋ น์ ๋ด๋ฆฌ๋ ๋ฐฉ์
- ์ฅ์ :
- ์ ์ฒด ํ์ดํ๋ผ์ธ์ ์ํฌํ๋ก์ฐ์ ์ํ๋ฅผ ์ค์(Web UI ๋ฑ)์์ ํ๋์ ๋ชจ๋ํฐ๋งํ๊ณ ๊ฐ์ํํ ์ ์์
- ์๋ฌ ๋ฐ์ ์ ์ค์์์ ์ฌ์๋๋ ๊ฒฐํจ ๊ฒฉ๋ฆฌ๋ฅผ ์ฆ๊ฐ ํต์ ํ ์ ์์
- ๋จ์ :
- ์ค์ ์ค์ผ์คํธ๋ ์ดํฐ ์์ง์ด ๋ง๋น๋๊ฑฐ๋ ๋ฉํ๋ฐ์ดํฐ DB๊ฐ ๋ค์ด๋๋ฉด
- ์ ์ฒด ์์คํ ํ์ดํ๋ผ์ธ์ด ๋ง๋น๋๋ ๋จ์ผ ์ฅ์ ์ (SPOF, Single Point of Failure)์ด ๋ ์ ์์
- ์ค์ ์ค์ผ์คํธ๋ ์ดํฐ ์์ง์ด ๋ง๋น๋๊ฑฐ๋ ๋ฉํ๋ฐ์ดํฐ DB๊ฐ ๋ค์ด๋๋ฉด
- ์ฝ๋ ์ค๊ทธ๋ํผ ํจํด (Choreography Pattern) - ์ด๋ฒคํธ ๋ถ์ฐํ
- ์ค์์ ํต์ ์ ์์ด, ๊ฐ ์๋น์ค๋ค์ด ๋ฉ์์ง ๋ธ๋ก์ปค(Kafka, RabbitMQ)๋ฅผ ํตํด ์ด๋ฒคํธ(Event)๋ฅผ ๋ฐํํ๊ณ ์์ ํ๋ฉฐ
-
์์จ์ ์ผ๋ก ์ถค์ถ๋ฏ ์ํธ์์ฉํ๋ ๋ฌด์ฉ(Choreography) ๋ฐฉ์
- ์ฅ์ :
- ์๋น์ค ๊ฐ์ ๊ฒฐํฉ๋(Coupling)๊ฐ ๊ทน๋๋ก ๋ฎ์ผ๋ฉฐ,
- ํน์ ์๋น์ค๊ฐ ์ฃฝ์ด๋ ๋ค๋ฅธ ์๋น์ค๋ ์ด๋ฒคํธ๋ฅผ ๊ณ์ ์ฒ๋ฆฌํ ์ ์์ด ํ์ฅ์ฑ๊ณผ ๊ฐ์ฉ์ฑ์ด ๋ฐ์ด๋จ
- ๋จ์ :
- ํ์ดํ๋ผ์ธ์ ์ ์ฒด ๋ฐ์ดํฐ ํ๋ฆ์ ํ๋์ ํ์ ํ๊ธฐ ์ด๋ ต๊ณ ,
- ํน์ ๊ตฌ๊ฐ์์ ๋ฐ์ดํฐ ์ ํฉ์ฑ์ด ๊นจ์ก์ ๋
- ์ญ์ถ์ (๋๋ฒ๊น ) ๋ฐ ๋ถ์ฐ ํธ๋์ญ์ ๋กค๋ฐฑ(Saga ํจํด ๊ตฌํ ๋ฑ)์ ๋์ด๋๊ฐ ๋น์ฝ์ ์ผ๋ก ์์น
- AI ๋ฐ ๋ฐ์ดํฐ ํ์ดํ๋ผ์ธ์ ์ ํ:
- ๋ฐ์ดํฐ ์์ง ๐กฒ ์ ์ฒ๋ฆฌ ๐กฒ ๋ชจ๋ธ ํ์ต์ผ๋ก ์ด์ด์ง๋ ์๊ฒฉํ ์ ํ ๊ด๊ณ์ ์ธ๊ณผ๊ด๊ณ๊ฐ ์ค์ํ
- AI/๋ฐ์ดํฐ ์์ง๋์ด๋ง ์์ญ์์๋ ๊ฐ์์ฑ๊ณผ ์ ์ด๊ถ์ด ๋ช ํํ โ์ค์ผ์คํธ๋ ์ด์ ํจํด(Apache Airflow ๋ฑ)โ์ ์๋์ ์ผ๋ก ์ ํธ
-
1.2 ์ค์ผ์คํธ๋ ์ด์ ์ํคํ ์ฒ์ 4๋ ํต์ฌ ์ปดํฌ๋ํธ
-
ํ๋์ ์ธ ์ค์ผ์คํธ๋ ์ด์ ์์ง์ ๋ด๋ถ์ ์ผ๋ก ๊ณ ๋์ ๋ถ์ฐ ์์คํ ๊ตฌ์กฐ๋ฅผ ์ฑํํ๊ณ ์์
- ์ปจํธ๋กค ํ๋ ์ธ & ์ค์ผ์ค๋ฌ (Control Plane & Scheduler):
- ์ ์ฒด ํ์ดํ๋ผ์ธ์ ๋ผ๋(DAG)๋ฅผ ํด์ํ๊ณ ,
- ๊ฐ ํ์คํฌ์ ์์กด์ฑ๊ณผ ์ง์ ์ฐจ์(In-degree)๋ฅผ ๊ณ์ฐํ์ฌ
- ์คํ ๊ฐ๋ฅํ ์์ ์ ์ ๋ณํ๋ ์ญํ
- ๋ฉํ๋ฐ์ดํฐ ์ ์ฅ์ (Metadata Repository):
- ๋ชจ๋ ์ํฌํ๋ก์ฐ์ ์คํ ์ด๋ ฅ, ํ์คํฌ์ ๋ผ์ดํ์ฌ์ดํด ์ํ(Queued, Running, Failed, Success), ์ ์ญ ์ค์ ๊ฐ ๋ฑ์ ์๊ตฌ ์ ์ฅํ๋ ์ํคํ ์ฒ์ ์ฌ์ฅ๋ถ
- RDBMS ์ฃผ๋ก ์ฌ์ฉ
- ์คํ๊ธฐ ๋ฐ ํ ์ธํ๋ผ (Executor & Message Queue):
- ์ค์ผ์ค๋ฌ๊ฐ ์ ๋ณํ ํ์คํฌ๋ฅผ ์ค์ ์ผ๊พผ(Worker)๋ค์๊ฒ ์์ ํ๊ฒ ์ ๋ฌํ๊ธฐ ์ํ ๋ฒํผ ๋ฐ ์ค๊ณ ๊ณ์ธต
- Redis, RabbitMQ ๊ฐ์ ๋ธ๋ก์ปค๋ ์ฟ ๋ฒ๋คํฐ์ค API API์ ์ฐ๋๋จ
- ๋ถ์ฐ ์์ปค ํด๋ฌ์คํฐ (Distributed Workers):
- ์ค์ผ์คํธ๋ ์ดํฐ์ ๋ช ๋ น์ ๋ฐ์ ์ค์ ์ปดํจํ ์ฐ์ฐ(API ํธ์ถ, SQL ํธ๋ฆฌ๊ฑฐ, Spark ์ง๋ ฌํ)์ ์ํํ๋ ๋ฌผ๋ฆฌ/๋ ผ๋ฆฌ์ ๋ ธ๋
- ์ปจํธ๋กค ํ๋ ์ธ & ์ค์ผ์ค๋ฌ (Control Plane & Scheduler):
1.3 ์ค์ผ์คํธ๋ ์ด์ ์ค๊ณ ์ ํ์ ์ํคํ ์ฒ ์์น
-
์ฑ๊ณต์ ์ธ ์ค์ผ์คํธ๋ ์ด์ ํ์ดํ๋ผ์ธ์ ๊ตฌ์ถํ๊ธฐ ์ํด ์ํคํ ์ฒ ๋ ๋ฒจ์์ ๋ฐ๋์ ์ค์ํด์ผ ํ๋ ์์ง๋์ด๋ง ์์น
- ์ปดํจํ
๊ณผ ์ค์ผ์คํธ๋ ์ด์
์ ๋ถ๋ฆฌ (Decoupling)
- ๊ฐ์ฅ ์ค์ํ ์์น ๐กฒ โ์ค์ผ์คํธ๋ ์ดํฐ๋ ์ ํธ๋ฑ ์ญํ ๋ง ํด์ผ์ง, ์ง์ ์ฐจ๊ฐ ๋์ด์๋ ์ ๋๋คโ๋ ๋ฒ์น
- ์ค์ผ์คํธ๋ ์ดํฐ ์์ปค ์์ฒด์ ๋ฉ๋ชจ๋ฆฌ์ CPU๋ฅผ ์๋ชจํ์ฌ ๋์ฉ๋ ๋ฐ์ดํฐ๋ฅผ ์ฒ๋ฆฌ(์: Pandas ๋ฐ์ดํฐ ๋ณํ)ํ๋ฉด ์์คํ ์ ์ฒด๊ฐ ๋ง๋น๋จ
- ๋ฌด๊ฑฐ์ด ์ฐ์ฐ์ ๋ถ์ฐ ์ปดํจํ ์์ง(Spark, Ray, Trino)์ด๋ ์ธ๋ถ DB ์์ง์ ์์ํ๊ณ ,
- ์ค์ผ์คํธ๋ ์ดํฐ๋ ์คํ ๋ช ๋ น(Trigger)๊ณผ ์๋ฃ ์ฌ๋ถ ํ์ธ(Polling)๋ง ์ํํด์ผ ํจ
- ๋ฉฑ๋ฑ์ฑ (Idempotency) ์ธํ๋ผ ๊ตฌ์ถ
- ์ค์ผ์คํธ๋ ์ด์
์ํคํ
์ฒ์์๋ ํน์ ํ์คํฌ๊ฐ ์คํจํ์ฌ ์ฌ์คํ(Retry)๋๊ฑฐ๋ ๊ณผ๊ฑฐ ํน์ ์์ ์ผ๋ก ๋์๊ฐ ๋ฐฑํ(Backfill)์ ์ํํ ๋,
- ๋ช ๋ฒ์ ๋ค์ ์คํํด๋ ํ๊ฒ ์ ์ฅ์์ ์ต์ข ๋ฐ์ดํฐ์ ๊ฒฐ๊ณผ๊ฐ ํญ์ ๋์ผํจ์ ๋ณด์ฅํด์ผ ํจ
- ๋ฐ์ดํฐ ์์ค๋ฅผ ๊ฒฉ๋ฆฌํ ์ ์๋ ๋ ผ๋ฆฌ์ ์์ ๋ณ์(Logical Date) ๋ฐ์ธ๋ฉ ๋ฐ ์ ์ฅ์์ Upsert ๋ฉ์ปค๋์ฆ ์ค๊ณ๊ฐ ์ํคํ ์ฒ์ ๋ด์ฌ๋์ด์ผ ํจ
- ์ค์ผ์คํธ๋ ์ด์
์ํคํ
์ฒ์์๋ ํน์ ํ์คํฌ๊ฐ ์คํจํ์ฌ ์ฌ์คํ(Retry)๋๊ฑฐ๋ ๊ณผ๊ฑฐ ํน์ ์์ ์ผ๋ก ๋์๊ฐ ๋ฐฑํ(Backfill)์ ์ํํ ๋,
- ๊ฒฐํฉ ๊ฒฉ๋ฆฌ ๋ฐ ์์ ๊ฒฉ๋ฆฌ (Isolation)
- ํ์ดํ๋ผ์ธ ๋ด์ ๊ฐ ๋จ๊ณ๋ ์ํธ ๊ฐ์ ๋ผ์ด๋ธ๋ฌ๋ฆฌ ์์กด์ฑ์ด๋ ํ๋์จ์ด ์์ ์๋น์ ์ํฅ์ ์ฃผ์ง ์์์ผ ํจ
- AI ํ์ดํ๋ผ์ธ์์๋ ์ผ๋ฐ ๊ฐ๋ฒผ์ด SQL ์ ์ฒ๋ฆฌ ํ์คํฌ์ ๋์ฉ๋ GPU ๊ฐ์์ด ํ์ํ ๋ชจ๋ธ ํ์ต ํ์คํฌ๊ฐ ๊ณต์กดํ๋ฏ๋ก,
- ํ์คํฌ ๋จ์๋ฅผ ์ปจํ ์ด๋(Docker Pod) ํํ๋ก ๋์ ๊ฒฉ๋ฆฌํ์ฌ
- ํ์ํ ์ธํ๋ผ์ ์ค์๊ฐ ๋ฐฐ์ ํ๋ ์ํคํ
์ฒ(์:
KubernetesPodOperator)๋ฅผ ๊ตฌ์ถํ๋ ๊ฒ์ด ์ต์
- ์ปดํจํ
๊ณผ ์ค์ผ์คํธ๋ ์ด์
์ ๋ถ๋ฆฌ (Decoupling)
1.4 ๊ธฐ์ ์คํ๋ณ ์ญํ ์ ์
- Ingestion (์์ง):
- ๋๋ด์ธ ์์ค(API, DB, ์นํ )๋ก๋ถํฐ ๊ฐ์ด๋, ๋งค๋ด์ผ ๋ฑ์ ์์ ๋ฐ์ดํฐ๋ฅผ ์์ง
- Data Lake (์ ์ฅ):
- ๋น์ฉ์ด ์ ๋ ดํ๊ณ ํ์ฅ์ฑ์ด ๋ฐ์ด๋ ์ค๋ธ์ ํธ ์คํ ๋ฆฌ์ง(์คํ์์ค MinIO ๋๋ AWS S3)๋ฅผ ์ํคํ ์ฒ์ ์ค์ฌ์ ๋
- Distributed Processing (Apache Spark):
- ๋ฐ์ดํฐ ๋ถ์ ๋ฐ ๋๊ท๋ชจ ๋ถ์ฐ ์ฐ์ฐ์ ํ์ค ์์ง
- ์์ง๋ ๋์ฉ๋ ํ ์คํธ์ ์ ์ , ํํ์ ๋ถ์, ํ ํฐ ํฌ๊ธฐ ๊ธฐ๋ฐ ๋ถ์ฐ ์ฒญํน(Chunking)์ ์ฒ๋ฆฌํจ
- Vector Database (Qdrant):
- ๊ณ ์ฐจ์ ๋ฒกํฐ ์๋ฒ ๋ฉ ๋ฐ์ดํฐ๋ฅผ ์ ์ฅํ๊ณ ,
- ์ฝ์ฌ์ธ ์ ์ฌ๋(Cosine Similarity) ๋ฑ์ ์๊ณ ๋ฆฌ์ฆ์ ๊ธฐ๋ฐ์ผ๋ก ๋ฐ์ง ๊ฒ์(Dense Retrieval)์ ์ด๊ณ ์์ผ๋ก ์ํํ๋
- ํ์ด๋ธ๋ฆฌ๋ ๊ฒ์ ์์ง
1.5 ์ค์ผ์คํธ๋ ์ด์ ํ๋ก์ธ์ค ํ์๋ผ์ธ
- Airflow ์ค์ผ์ค๋ฌ๊ฐ ์ ์ฒด ํ์ดํ๋ผ์ธ์ ์๋ช ์ฃผ๊ธฐ๋ฅผ ๊ด๋ฆฌํ๋ 4๋จ๊ณ ํ๋ก์ธ์ค
- [Task 1] Ingestion & Lake Landing (์์ง ๋ฐ ์ ์ฌ):
- ์ธ๋ถ ๋ฐ์ดํฐ๋ฅผ ๋ค์ด๋ก๋ํ์ฌ
- MinIO์
raw-data/๋ฒํท์ ์ ์ฅ
- [Task 2] Spark ๋ถ์ฐ ์ฐ์ฐ ํธ๋ฆฌ๊ฑฐ (Spark-Submit):
- Airflow๊ฐ Spark Operator๋ฅผ ํตํด ์ ์ฒ๋ฆฌ ์์ ์ ๋ช ๋ น
- Spark ํด๋ฌ์คํฐ๊ฐ MinIO์ ์์ ๋ฐ์ดํฐ๋ฅผ ์ฝ์ด์ ์ฒญํน(Chunking)์ ์ํ
processed-data/๋ฒํท์ Parquet ํ์์ผ๋ก ์ ์ฅ
- [Task 3] ๋ณ๋ ฌ ๋ถ์ฐ ์๋ฒ ๋ฉ ๋ฐ Vector DB Upsert:
- ๊ฐ๊ณต๋ ํ ์คํธ ์ฒญํฌ๋ค์ ์ฝ์ด์
- ๋ก์ปฌ AI ์๋ฒ ๋ฉ ๋ชจ๋ธ(Ollama/HuggingFace)์ ํตํด ๊ณ ์ฐจ์ ๋ฒกํฐ๋ก ๋ณํํ ๋ค,
- Qdrant์ HNSW ๊ทธ๋ํ ์ธ๋ฑ์ค์ Upsert
- [Task 4] ์ธ๋ฑ์ค ์ ๋น ๋ฐ ์บ์ ํด๋ฆฐ์
:
- ๋ฉํ๋ฐ์ดํฐ ๊ฐฑ์ ๋ฐ ๋ฆฌ์์ค ํด์
2. ์ข ํฉ ์ค์ต ์์ ์ฝ๋ 1
- ํ๋์ ์ธ AI ์ธํ๋ผ์ ํ์ค ๊ตฌ์กฐ์ธ ํ์ด๋ธ๋ฆฌ๋ RAG(๊ฒ์ ์ฆ๊ฐ ์์ฑ) ํ๋ซํผ์ ๋ฐ์ดํฐ ๊ณต๊ธ์ ์ ์ํคํ ์ฒ ๊ด์ ์์ ํ๋กํ ํ์ดํํ ํต์ฌ ์ค๊ณ๋
- โ์์ง ๐กฒ Data Lake ๐กฒ Apache Spark ๐กฒ VectorDBโ ๊ตฌ์กฐ๋ฅผ ๋จ์ผ ํ์ผ๋ก ๊ตฌํํ Airflow DAG
- ์ค๋ฌด์์๋ Spark ์ ์ฒ๋ฆฌ ๋ก์ง์ ๋ณ๋์
*.pyํ์ผ๋ก ๋ถ๋ฆฌํ์ฌSparkSubmitOperator๋ก ํธ์ถ - ์ค์ต์์ ์์๋ ๊ฐ๋ ์ฑ์ ์ํด PySpark ์ ์ฒ๋ฆฌ ๋ฐ Qdrant ์ ์ฌ๋ฅผ ํตํฉ ๊ตฌํํจ
- ์ค๋ฌด์์๋ Spark ์ ์ฒ๋ฆฌ ๋ก์ง์ ๋ณ๋์
2.1 ์ํคํ ์ฒ์ ๊ตฌ์ฑ ์์ (Components)
- ์ค์ ์งํ์ (Apache Airflow DAG):
task_ingest์task_spark_and_vector๋ผ๋ ๋ ๊ฐ์ ์คํ ๋ ธ๋๋ฅผ ํต์ - ์ด๋ค ์์ ์ด ๋จผ์ ์คํ๋์ด์ผ ํ๋์ง ์ ํ ๊ด๊ณ๋ฅผ ๊ท์
- ์ฅ์ ๋ฐ์ ์ ์๋์ผ๋ก ์ฌ์๋(
retries: 2)ํ๋ ๊ด๋ฆฌ ๊ณ์ธต
- ์์ ๋ฐ์ดํฐ ๋ ์ดํฌ (MinIO):
- ์ค๋ธ์ ํธ ์คํ ๋ฆฌ์ง ์์ญ
- ๋น์ ํ ๋ฐ์ดํฐ(๋งค๋ด์ผ ํ
์คํธ)๋ฅผ ๊ฐ๊ณต๋์ง ์์ ์์ ์ค๋
์ท ์ํ ๊ทธ๋๋ก ์์ ํ๊ฒ ์๊ตฌ ์ ์ฅ(
factory-raw-logs๋ฒํท)
- ๋ถ์ฐ ์ปดํจํ
๊ฐ๊ณต ๋ฐ ๊ณ ์ฐจ์ ์ ์ฅ์ (Spark & Vector DB):
- ์ฝ๋๋ ์ธ๋ผ์ธ์ผ๋ก ๊ฐ๋ณ๊ฒ ๊ตฌํ๋์ด ์์ผ๋
- ๋ ผ๋ฆฌ์ ์ผ๋ก๋ ๋ํ ํ ์คํธ ์ ์ ์ฐ์ฐ(Spark)์ ์ํํ์ฌ
- ์ด๋ฅผ ๊ณ ์ ๊ทธ๋ํ ํ์ ์ธ๋ฑ์ค์ธ HNSW ๊ทธ๋ํ ๊ตฌ์กฐ(Qdrant)๋ก ๋๊ธฐํํ๋
- AI ๋ฐ์ดํฐ ์๋น ์ปดํฌ๋ํธ
2.2 ์ฝ๋์ ์ ์ ๊ตฌ์กฐ (Static Structure)
-
์ฝ๋๋ ํฌ๊ฒ ์ธ ๊ฐ์ ์์ญ(Configuration, Implementation, Orchestration)์ผ๋ก ๋ ์ด์ด๊ฐ ๋๋์ด ์์
[ ๊ตฌ์กฐ์ ๋ ์ด์ด ๊ตฌ์ฑ ] โโ 1. ์ ์ญ ์ธํ๋ผ ํ ํด๋ก์ง ์ ์ ๋ ์ด์ด (์ค์ ๊ตฌ์ญ: MINIO_URL, QDRANT_URL ๋ฑ) โโ 2. ๋น์ฆ๋์ค ๋ก์ง ๊ตฌํ ๋ ์ด์ด (์คํ ํจ์: fn_ingest_to_lake, fn_spark_...) โโ 3. ์ํฌํ๋ก์ฐ ์ค์ผ์คํธ๋ ์ด์ ๋ ์ด์ด (with DAG(...) ๊ตฌ๋ฌธ ๋ฐ ์์กด์ฑ ์ ์ธ) - ์ ์ญ ์ธํ๋ผ ํ ํด๋ก์ง ์ ์ ๋ ์ด์ด (์ต์๋จ ๊ตฌ์ญ)
- ์ธ๋ถ ์ธํ๋ผ ์์คํ
๋ค์ ์ ์ ์ฃผ์์ ์ธ์ฆ ์ ๋ณด(
MINIO_ACCESS,QDRANT_URL๋ฑ) ๋ฐ ์ ์ฌ ๋์์ด ๋ ๋ฌผ๋ฆฌ ๊ณต๊ฐ(BUCKET_NAME,COLLECTION_NAME)์ ์ ์ญ ๋ณ์๋ก ๊ท์
- ์ธ๋ถ ์ธํ๋ผ ์์คํ
๋ค์ ์ ์ ์ฃผ์์ ์ธ์ฆ ์ ๋ณด(
- ๋น์ฆ๋์ค ๋ก์ง ๊ตฌํ ๋ ์ด์ด (์ค๋จ ๊ตฌ์ญ)
fn_ingest_to_lake():- MinIO SDK ํด๋ผ์ด์ธํธ๋ฅผ ์ ์ธํ๊ณ
bucket_existsAPI๋ฅผ ํตํด ๋ฐฉ์ด์ ์ฝ๋(Defensive Coding)๋ฅผ ๊ตฌ์ถํ ๋ค,- ์ค๋งํธํฉํ ๋ฆฌ ๋๋ฉ์ธ์ ์ ๋๊ณ ์ ๋น ๋งค๋ด์ผ ๋น์ ํ ๋ฐ์ดํฐ๋ฅผ
- ๋ ์ดํฌ์ ์ ์ฌํ๋ ํจ์
fn_spark_processing_and_vector_upsert():- ๋ฐ์ดํฐ ๋ ์ดํฌ์์ ํ์ผ์ ๋ค์ ์คํธ๋ฆฌ๋ฐ์ผ๋ก ์ฝ์ด์
replace์ฐ์ฐ์ผ๋ก ๋ ธ์ด์ฆ๋ฅผ ์ ๊ฑฐํ๊ณ ,- 50๊ธ์ ๋จ์๋ก ์๋ผ๋ด๋ ์๋ฏธ๋ก ์ ์ฒญํน(Chunking)์ ์ฒ๋ฆฌ
- ์ดํ Qdrant์ ANN(๊ทผ์ฌ ์ต๊ทผ์ ์ด์) ์ ์ฌ๋ ๊ฒ์์ ์ํด
- 5์ฐจ์ ๊ณต๊ฐ ๋ฒกํฐ ๊ตฌ์กฐ(
PointStruct)๋ก ํฌ๋งทํ ํ์ฌ ์ ์ฌ๋ฅผ ์ ๋ด
- ์ํฌํ๋ก์ฐ ์ค์ผ์คํธ๋ ์ด์
๋ ์ด์ด (ํ๋จ ๊ตฌ์ญ)
with DAG(...) as dag:- ์ปจํ ์คํธ ๋งค๋์ ๋ฅผ ํตํด Airflow ์ปดํ์ผ๋ฌ ๋ด๋ถ๋ก ์ง์
PythonOperator๋ค์ ํ์ฉํด- ์์ ๊ตฌํํ ๋น์ฆ๋์ค ํ์ด์ฌ ํจ์๋ค์
- Airflow ๊ฐ์์ ์ธ ํ์คํฌ ๋ ธ๋๋ก ์ธ์คํด์คํ
2.3 ๋์ ๋ฐ์ดํฐ ํ๋ฆ ๋ฐ ๋ฉ์ปค๋์ฆ (Runtime Flow)
-
DAG๊ฐ ํธ๋ฆฌ๊ฑฐ๋๋ ์๊ฐ, ๋ฐ์ดํฐ์ ์ ์ด๊ถ์ ์ ํ ๊ด๊ณ์ ๋ง์ถฐ ๋ฌผ๋ฆฌ ์์์ ์ด๋ํ๊ฒ ๋จ
[ ๋ฐ์ดํฐ ๋ฐ ์ ์ด๊ถ ํ๋ฆ ํ์๋ผ์ธ ] (์ธ๋ถ ์์ค ๋ฐ์ดํฐ) โ โผ [ 1๋จ๊ณ: task_ingest ] โโโโ MinIO (enterprise-knowledge-lake) ์ ์ฅ ์๋ฃ โ โโ (์์กด์ฑ ์ ์ด๊ถ ์ดํ: '>>') โผ [ 2๋จ๊ณ: task_spark_and_vector ] โโ โ ๋ฐ์ดํฐ ๋ ์ดํฌ๋ก๋ถํฐ ์์ ํ ์คํธ ์คํธ๋ฆฌ๋ฐ Load โโ โก ๋ถ์ฐ ๋ฉ๋ชจ๋ฆฌ ๊ณต๊ฐ ์ฒญํน ์ฐ์ฐ (40~50์ ๋ถํ ) โโ โข Qdrant ๊ณ ์ฐจ์ ๋ฒกํฐ ํน์ง ๊ณต๊ฐ ๋งตํ (HNSW ์ธ๋ฑ์ฑ)- ์ง์
๊ด๋ฌธ ๋ฐ ๋ฐ์ดํฐ ๋ ์ดํฌ ์์ฐฉ (
task_ingest):- Airflow ์ค์ผ์ค๋ฌ๊ฐ ์ง์
์ฐจ์(
In-degree=0)๊ฐ ์ ๋ก์ธtask_ingest๋ฅผ ๊ฐ์ฅ ๋จผ์ ํ์ ๋ฃ๊ณ ์์ปค์ ๋ฐฐ์ - ๊ฐ์์ ๋น์ ํ ์์ค ํ
์คํธ ๋ฐ์ดํฐ๊ฐ
- ๋ฐ์ดํธ ์คํธ๋ฆผ(
io.BytesIO) ํํ๋ก ๋ณํ๋์ด - ๋คํธ์ํฌ ๋ง์ ํ๊ณ
- MinIO ์ค๋ธ์ ํธ ์คํ ๋ฆฌ์ง ๋ด๋ถ์
raw/manual_01.txt๊ฒฝ๋ก๋ก ์ ๋ก๋ ๋ฐ ๊ฒฉ๋ฆฌ
- ๋ฐ์ดํธ ์คํธ๋ฆผ(
- Airflow ์ค์ผ์ค๋ฌ๊ฐ ์ง์
์ฐจ์(
- ์์ฐจ ์ ์ด๊ถ ์ดํ (
>>):- 1๋จ๊ณ ํ์คํฌ๊ฐ ์ฑ๊ณต(
Success)์ผ๋ก ๋งํน๋๋ฉด, - ์์กด์ฑ ๊ฒฐํฉ ์ฐ์ฐ์(
>>)๋ฅผ ํ๊ณ - ์ ์ด๊ถ์ด ๋ค์ ํ์คํฌ์ธ
task_spark_and_vector๋ก ์์ ํ๊ฒ ์ ์ด
- 1๋จ๊ณ ํ์คํฌ๊ฐ ์ฑ๊ณต(
- ๋ฐ์ดํฐํ์ฐ์ค ๋ก๋ ๋ฐ ๋ฉ๋ชจ๋ฆฌ ์ฒญํน ์ฐ์ฐ:
- ๋ ๋ฒ์งธ ํ์คํฌ๊ฐ ๊ธฐ๋๋๋ฉด์ ๋ฐฉ๊ธ MinIO์ ๋ฐฑ์ ๋์๋ ์์ ํ์ผ์ ์ค๋ ์ท์ ๋ฉ๋ชจ๋ฆฌ๋ก ๋ค์ ๋ค์ด๋ก๋
replace๊ฐ๊ณต์ ํตํด ๋ฌธ์์ด ๋ ธ์ด์ฆ๋ฅผ ์ ์ ํ๊ณ ,- ์ฌ๋ผ์ด๋ ์๋์ฐ ๋ฐฉ์์ผ๋ก 50๊ธ์์ฉ ์ชผ๊ฐ์ง ํ
์คํธ ์คํธ๋ง ๋ฆฌ์คํธ(
chunks)๋ฅผ ์์ฑํ์ฌ - ๋ฉ๋ชจ๋ฆฌ์ ๋ถ์ฐ ๋ฐฐ์น
- ์ฌ๋ผ์ด๋ ์๋์ฐ ๋ฐฉ์์ผ๋ก 50๊ธ์์ฉ ์ชผ๊ฐ์ง ํ
์คํธ ์คํธ๋ง ๋ฆฌ์คํธ(
- ๋ฉฑ๋ฑ์ฑ ๊ธฐ๋ฐ ๊ณ ์ฐจ์ ๋ฒกํฐ์คํ ์ด ์ต์ข
๋๊ธฐํ:
- Qdrant ํด๋ผ์ด์ธํธ๊ฐ ๊ฐ๋๋์ด
- ๋ฒกํฐ์คํ ์ด ๋ด๋ถ์
Distance.COSINE๊ฑฐ๋ฆฌ๋ฅผ ์ฐ์ฐํ ์ ์๋ ์ธ๋ฑ์ค ๋ ์ด์ด๋ฅผ ์ ์ธ
- ๋ฒกํฐ์คํ ์ด ๋ด๋ถ์
- ์ฒญํน๋ ๋ฌธ์์ด ๋ฐ์ดํฐ ๊ฐ๊ฐ์ ์ ์ผํ ID ๊ณ ์ ๋ฒํธ(
idx + 100)๋ฅผ ๊ฐ์ ๋ก ๋ถ์ฌ - ์ด๋ ๊ฒ ๊ณ ์ ID๋ฅผ ๋ฐ์ธ๋ฉํจ์ผ๋ก์จ,
- ์ด ํ์ดํ๋ผ์ธ์ ํ๋ฃจ์ ์์ญ ๋ฒ ์ค๋ณตํด์ ๋ค์ ์คํํ๋๋ผ๋
- Qdrant ๋ด๋ถ์ ๋ฐ์ดํฐ๊ฐ ๋์ ๋์ง ์๊ณ ๋ฎ์ด์จ์ง๊ฒ(Upsert) ๋ง๋ค์ด
- ๋ฐ์ดํฐ์ ๋ฉฑ๋ฑ์ฑ(Idempotency)์ ์ต์ข ์์ํ๊ณ
- ํ์ดํ๋ผ์ธ ์ ์ฒด๊ฐ ์ ์ ์ข ๋ฃ๋จ
- Qdrant ํด๋ผ์ด์ธํธ๊ฐ ๊ฐ๋๋์ด
- ์ง์
๊ด๋ฌธ ๋ฐ ๋ฐ์ดํฐ ๋ ์ดํฌ ์์ฐฉ (
2.4 ์์ ์ฝ๋
docker-compose.yml- ๋์ปค ์ปจํ
์ด๋ ํ๊ฒฝ ์ค์
- Apache Airflow๋ ๊ณต์ ์ฌ์ดํธ์์ ์ ๊ณตํ๋ ์ต์ ํ์ผ์ ๋ค์ด๋ก๋ํด์ ์ด์ฉํ ๊ฒ
- MinIO, Qdrant์ ์ค์ ์ ์ถ๊ฐํ ๊ฒ
x-airflow-common: &airflow-common # In order to add custom dependencies or upgrade provider distributions you can use your extended image. # Comment the image line, place your Dockerfile in the directory where you placed the docker-compose.yaml # and uncomment the "build" line below, Then run `docker-compose build` to build the images. image: ${AIRFLOW_IMAGE_NAME:-apache/airflow:3.3.0} # build: . env_file: - ${ENV_FILE_PATH:-.env} environment: &airflow-common-env AIRFLOW__CORE__EXECUTOR: CeleryExecutor AIRFLOW__CORE__AUTH_MANAGER: airflow.providers.fab.auth_manager.fab_auth_manager.FabAuthManager AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@postgres/airflow AIRFLOW__CELERY__RESULT_BACKEND: db+postgresql+psycopg2://airflow:airflow@postgres/airflow AIRFLOW__CELERY__BROKER_URL: redis://:@redis:6379/0 AIRFLOW__CORE__FERNET_KEY: ${FERNET_KEY} AIRFLOW__CORE__DAGS_ARE_PAUSED_AT_CREATION: 'true' AIRFLOW__CORE__LOAD_EXAMPLES: 'true' AIRFLOW__CORE__EXECUTION_API_SERVER_URL: 'http://airflow-apiserver:8080/execution/' AIRFLOW__API_AUTH__JWT_SECRET: ${AIRFLOW__API_AUTH__JWT_SECRET:-airflow_jwt_secret} AIRFLOW__API_AUTH__JWT_ISSUER: ${AIRFLOW__API_AUTH__JWT_ISSUER:-airflow} # yamllint disable rule:line-length # Use simple http server on scheduler for health checks # See https://airflow.apache.org/docs/apache-airflow/stable/administration-and-deployment/logging-monitoring/check-health.html#scheduler-health-check-server # yamllint enable rule:line-length AIRFLOW__SCHEDULER__ENABLE_HEALTH_CHECK: 'true' # WARNING: Use _PIP_ADDITIONAL_REQUIREMENTS option ONLY for a quick checks # for other purpose (development, test and especially production usage) build/extend Airflow image. _PIP_ADDITIONAL_REQUIREMENTS: ${_PIP_ADDITIONAL_REQUIREMENTS:-} # The following line can be used to set a custom config file, stored in the local config folder AIRFLOW_CONFIG: '/opt/airflow/config/airflow.cfg' volumes: - ${AIRFLOW_PROJ_DIR:-.}/dags:/opt/airflow/dags - ${AIRFLOW_PROJ_DIR:-.}/logs:/opt/airflow/logs - ${AIRFLOW_PROJ_DIR:-.}/config:/opt/airflow/config - ${AIRFLOW_PROJ_DIR:-.}/plugins:/opt/airflow/plugins user: "${AIRFLOW_UID:-50000}:0" depends_on: &airflow-common-depends-on redis: condition: service_healthy postgres: condition: service_healthy services: postgres: image: postgres:16 environment: POSTGRES_USER: airflow POSTGRES_PASSWORD: airflow POSTGRES_DB: airflow volumes: - postgres-db-volume:/var/lib/postgresql/data healthcheck: test: ["CMD", "pg_isready", "-U", "airflow"] interval: 10s retries: 5 start_period: 5s restart: always redis: # Redis is limited to 7.2-bookworm due to licencing change # https://redis.io/blog/redis-adopts-dual-source-available-licensing/ image: redis:7.2-bookworm expose: - 6379 healthcheck: test: ["CMD", "redis-cli", "ping"] interval: 10s timeout: 30s retries: 50 start_period: 30s restart: always airflow-apiserver: <<: *airflow-common command: api-server ports: - "8080:8080" healthcheck: test: ["CMD", "curl", "--fail", "http://localhost:8080/api/v2/monitor/health"] interval: 30s timeout: 10s retries: 5 start_period: 30s restart: always depends_on: <<: *airflow-common-depends-on airflow-init: condition: service_completed_successfully airflow-scheduler: <<: *airflow-common command: scheduler healthcheck: test: ["CMD-SHELL", 'airflow jobs check --job-type SchedulerJob --hostname "$${HOSTNAME}"'] interval: 30s timeout: 10s retries: 5 start_period: 30s restart: always depends_on: <<: *airflow-common-depends-on airflow-init: condition: service_completed_successfully airflow-dag-processor: <<: *airflow-common command: dag-processor healthcheck: test: ["CMD-SHELL", 'airflow jobs check --job-type DagProcessorJob --hostname "$${HOSTNAME}"'] interval: 30s timeout: 10s retries: 5 start_period: 30s restart: always depends_on: <<: *airflow-common-depends-on airflow-init: condition: service_completed_successfully airflow-worker: <<: *airflow-common command: celery worker healthcheck: # yamllint disable rule:line-length test: ["CMD-SHELL", 'celery --app airflow.providers.celery.executors.celery_executor.app inspect ping -d "celery@$${HOSTNAME}" || celery --app airflow.executors.celery_executor.app inspect ping -d "celery@$${HOSTNAME}"'] interval: 30s timeout: 10s retries: 5 start_period: 30s environment: <<: *airflow-common-env # Required to handle warm shutdown of the celery workers properly # See https://airflow.apache.org/docs/docker-stack/entrypoint.html#signal-propagation DUMB_INIT_SETSID: "0" restart: always depends_on: <<: *airflow-common-depends-on airflow-apiserver: condition: service_healthy airflow-init: condition: service_completed_successfully airflow-triggerer: <<: *airflow-common command: triggerer healthcheck: test: ["CMD-SHELL", 'airflow jobs check --job-type TriggererJob --hostname "$${HOSTNAME}"'] interval: 30s timeout: 10s retries: 5 start_period: 30s restart: always depends_on: <<: *airflow-common-depends-on airflow-init: condition: service_completed_successfully airflow-init: <<: *airflow-common entrypoint: /bin/bash # yamllint disable rule:line-length command: - -c - | if [[ -z "${AIRFLOW_UID}" ]]; then echo echo -e "\033[1;33mWARNING!!!: AIRFLOW_UID not set!\e[0m" echo "If you are on Linux, you SHOULD follow the instructions below to set " echo "AIRFLOW_UID environment variable, otherwise files will be owned by root." echo "For other operating systems you can get rid of the warning with manually created .env file:" echo " See: https://airflow.apache.org/docs/apache-airflow/stable/howto/docker-compose/index.html#setting-the-right-airflow-user" echo export AIRFLOW_UID=$$(id -u) fi one_meg=1048576 mem_available=$$(($$(getconf _PHYS_PAGES) * $$(getconf PAGE_SIZE) / one_meg)) cpus_available=$$(grep -cE 'cpu[0-9]+' /proc/stat) disk_available=$$(df / | tail -1 | awk '{print $$4}') warning_resources="false" if (( mem_available < 4000 )) ; then echo echo -e "\033[1;33mWARNING!!!: Not enough memory available for Docker.\e[0m" echo "At least 4GB of memory required. You have $$(numfmt --to iec $$((mem_available * one_meg)))" echo warning_resources="true" fi if (( cpus_available < 2 )); then echo echo -e "\033[1;33mWARNING!!!: Not enough CPUS available for Docker.\e[0m" echo "At least 2 CPUs recommended. You have $${cpus_available}" echo warning_resources="true" fi if (( disk_available < one_meg * 10 )); then echo echo -e "\033[1;33mWARNING!!!: Not enough Disk space available for Docker.\e[0m" echo "At least 10 GBs recommended. You have $$(numfmt --to iec $$((disk_available * 1024 )))" echo warning_resources="true" fi if [[ $${warning_resources} == "true" ]]; then echo echo -e "\033[1;33mWARNING!!!: You have not enough resources to run Airflow (see above)!\e[0m" echo "Please follow the instructions to increase amount of resources available:" echo " https://airflow.apache.org/docs/apache-airflow/stable/howto/docker-compose/index.html#before-you-begin" echo fi echo echo "Creating missing opt dirs if missing:" echo mkdir -v -p /opt/airflow/{logs,dags,plugins,config} echo echo "Airflow version:" /entrypoint airflow version echo echo "Files in shared volumes:" echo ls -la /opt/airflow/{logs,dags,plugins,config} echo echo "Running airflow config list to create default config file if missing." echo /entrypoint airflow config list >/dev/null echo echo "Files in shared volumes:" echo ls -la /opt/airflow/{logs,dags,plugins,config} echo echo "Change ownership of files in /opt/airflow to ${AIRFLOW_UID:-50000}:0" echo chown -R "${AIRFLOW_UID:-50000}:0" /opt/airflow/ echo echo "Change ownership of files in shared volumes to ${AIRFLOW_UID:-50000}:0" echo chown -v -R "${AIRFLOW_UID:-50000}:0" /opt/airflow/{logs,dags,plugins,config} echo echo "Files in shared volumes:" echo ls -la /opt/airflow/{logs,dags,plugins,config} # yamllint enable rule:line-length environment: <<: *airflow-common-env _AIRFLOW_DB_MIGRATE: 'true' _AIRFLOW_WWW_USER_CREATE: 'true' _AIRFLOW_WWW_USER_USERNAME: ${_AIRFLOW_WWW_USER_USERNAME:-airflow} _AIRFLOW_WWW_USER_PASSWORD: ${_AIRFLOW_WWW_USER_PASSWORD:-airflow} _PIP_ADDITIONAL_REQUIREMENTS: '' user: "0:0" airflow-cli: <<: *airflow-common profiles: - debug environment: <<: *airflow-common-env CONNECTION_CHECK_MAX_COUNT: "0" # Workaround for entrypoint issue. See: https://github.com/apache/airflow/issues/16252 command: - bash - -c - airflow depends_on: <<: *airflow-common-depends-on # You can enable flower by adding "--profile flower" option e.g. docker-compose --profile flower up # or by explicitly targeted on the command line e.g. docker-compose up flower. # See: https://docs.docker.com/compose/profiles/ flower: <<: *airflow-common command: celery flower profiles: - flower ports: - "5555:5555" healthcheck: test: ["CMD", "curl", "--fail", "http://localhost:5555/"] interval: 30s timeout: 10s retries: 5 start_period: 30s restart: always depends_on: <<: *airflow-common-depends-on airflow-init: condition: service_completed_successfully minio-local: image: minio/minio:RELEASE.2024-01-11T07-46-16Z container_name: minio-local ports: - "9000:9000" # ํ์ด์ฌ SDK๊ฐ ์ ์ํ API ํต์ ํฌํธ - "9001:9001" # ์น UI ์ฝ์ ์ด๋๋ฏผ ํฌํธ environment: MINIO_ROOT_USER: minioadmin MINIO_ROOT_PASSWORD: minioadminpassword volumes: - minio-data-volume:/data command: server /data --console-address ":9001" networks: - airflow-net qdrant-local: image: qdrant/qdrant:v1.18.0 container_name: qdrant-local ports: - "6333:6333" # REST API ๋ฐ ๋์๋ณด๋ ์ง์ ํฌํธ volumes: - qdrant-data-volume:/qdrant/storage networks: - airflow-net networks: airflow-net: name: local-ai-platform-net driver: bridge volumes: postgres-db-volume: minio-data-volume: qdrant-data-volume:
- ๋์ปค ์ปจํ
์ด๋ ํ๊ฒฝ ์ค์
-
.env์ค์ AIRFLOW_UID=1000 FERNET_KEY=************************************* _PIP_ADDITIONAL_REQUIREMENTS=minio qdrant-client- MinIO, Qdrant๋ฅผ ์ธ์ํ ์ ์๋๋ก ์ปจํ
์ด๋ ์์ _PIP_ADDITIONAL_REQUIREMENTS ์ค์ ์ ํตํด ๋ผ์ด๋ธ๋ฌ๋ฆฌ ์ค์น
- ์ปจํ
์ด๋์ MinIO, Qdrant์ ํ์ด์ฌ ๋ผ์ด๋ธ๋ฌ๋ฆฌ์ ๋ฒ์ ์ด ๋ค๋ฅผ ๊ฒฝ์ฐ์๋ ์ค๋ฅ๊ฐ ๋ฐ์ํ ์ ์์
- ์: ์ปจํ ์ด๋์ Qdrant: 1.12.0 / Airflow ์ปจํ ์ด๋์์ ์์ฒญ์ ์ํด ์ค์นํ Qdrant ๋ผ์ด๋ธ๋ฌ๋ฆฌ: 1.18.0
- ์ปจํ
์ด๋์ ๋ฒ์ ์ ์์ ํ๊ฑฐ๋ ํ์ด๋ธ๋ฌ๋ฆฌ ์ค์น ์์ฒญ์์ ๋ฒ์ ์ ์ง์ ํ ๊ฒ
- ์ปจํ
์ด๋์ MinIO, Qdrant์ ํ์ด์ฌ ๋ผ์ด๋ธ๋ฌ๋ฆฌ์ ๋ฒ์ ์ด ๋ค๋ฅผ ๊ฒฝ์ฐ์๋ ์ค๋ฅ๊ฐ ๋ฐ์ํ ์ ์์
-
์์ 2์์๋ Apache Airflow ์ด๋ฏธ์ง์ ์ปค์คํ ๋น๋ ์์ MinIO, Qdrant ๋ฑ์ ํฌํจ์ํค๊ณ
.env์์๋ ์ญ์ ํจ
๐กช ์ฑ๋ฅ ๋ฐ ๊ฐ๋ ์๊ฐ ์ ๊ฐ์ ํจ์ฌ ํจ์จ์ ์ - FERNET_KEY๊ฐ ์ค์ ๋์ง ์์ผ๋ฉด ์ง์์ ์ผ๋ก ๊ฒฝ๊ณ ๋ฌธ๊ตฌ๊ฐ ๋ฐ์ํจ
- docker-compose.yml ๋ด๋ถ์ ์ปจํ
์ด๋ ๊ฐ์๋งํผ ๋ฐ์
-
FERNET_KEY ์์ฑ ๋ฐฉ๋ฒ
python3 -c "from cryptography.fernet import Fernet; print(Fernet.generate_key().decode())
-
- docker-compose.yml ๋ด๋ถ์ ์ปจํ
์ด๋ ๊ฐ์๋งํผ ๋ฐ์
- MinIO, Qdrant๋ฅผ ์ธ์ํ ์ ์๋๋ก ์ปจํ
์ด๋ ์์ _PIP_ADDITIONAL_REQUIREMENTS ์ค์ ์ ํตํด ๋ผ์ด๋ธ๋ฌ๋ฆฌ ์ค์น
-
DAG ์์ฑ(advanced_ai_orchestration_pipeline.py)
#//file: "dags/advanced_ai_orchestration_pipeline.py" from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator from minio import Minio from qdrant_client import QdrantClient from qdrant_client.models import Distance, VectorParams, PointStruct import io, os # ์ธํ๋ผ ํ ํด๋ก์ง ์๋ํฌ์ธํธ ์ ์ MINIO_URL = "192.168.0.6:9000" MINIO_ACCESS = "minioadmin" MINIO_SECRET = "minioadminpassword" BUCKET_NAME = "enterprise-knowledge-lake" QDRANT_URL = "192.168.0.6:6333" COLLECTION_NAME = "factory_manual_vectors" # [Task 1] ์ธ๋ถ ์์ค๋ก๋ถํฐ ๋ฐ์ดํฐ๋ฅผ ์์งํ์ฌ ๋ฐ์ดํฐ ๋ ์ดํฌ์ ์ ์ฅ def fn_ingest_to_lake(): client = Minio(MINIO_URL, access_key=MINIO_ACCESS, secret_key=MINIO_SECRET, secure=False) if not client.bucket_exists(BUCKET_NAME): client.make_bucket(BUCKET_NAME) dummy_doc = "TIMESTAMP:2026-07-08 | MANUAL:์ค๋งํธํฉํ ๋ฆฌ ๊ฐ๋ณ ๋ฉด์ ์ ๋๊ณ ๊ณ์ธก๊ธฐ ์ฅ์ ์กฐ์น ๋งค๋ด์ผ. ์๋ ฅ ๊ฐํ ์ ๋ฒจ๋ธ ํธ์ค๋ฅผ ์ ๊ฒํ์ญ์์ค." client.put_object( bucket_name=BUCKET_NAME, object_name="raw/manual_01.txt", data=io.BytesIO(dummy_doc.encode('utf-8')), length=len(dummy_doc.encode('utf-8')), content_type="text/plain" ) print("[Success] ์์ ๋งค๋ด์ผ์ด Data Lake(MinIO)์ ์์ฐฉํ์ต๋๋ค.") # [Task 2 & 3] Apache Spark ์คํค๋ง๋ฅผ ๋ชจ๋ฐฉํ ํ ์คํธ ๊ฐ๊ณต ๋ฐ VectorDB ์ ์ฌ # (์ฃผ์: ์ค์ ๋ถ์ฐํ๊ฒฝ์์๋ PySpark๋ฅผ ํ์ฉํด ๋๋ ๊ฐ๊ณต ํ ๋ถ์ฐ Upsert ์ฒ๋ฆฌ) def fn_spark_processing_and_vector_upsert(): # 1. Lake์์ ๋ฐ์ดํฐ ์ฝ๊ธฐ (Spark์ ๋ฐ์ดํฐ ์์ค ๋ก๋ฉ ์ญํ ๋ชจ๋ฐฉ) minio_client = Minio(MINIO_URL, access_key=MINIO_ACCESS, secret_key=MINIO_SECRET, secure=False) response = minio_client.get_object(BUCKET_NAME, "raw/manual_01.txt") raw_text = response.read().decode('utf-8') response.close() # 2. Spark ๋น์ฆ๋์ค ๋ก์ง: ์๋ฏธ๋ก ์ ๊ฐ๊ณต ๋ฐ ์ฒญํน (Transformation) # ์ค์ ์ค๋ฌด์์๋ Spark DataFrame์ udf(User Defined Function)๋ฅผ ์ฌ์ฉํ์ฌ ๋ถ์ฐ ํด๋ฌ์คํฐ์์ ์ํ๋จ refined_text = raw_text.replace("TIMESTAMP:2026-07-08 | ", "[ํ์ธ ์๋ฃ] ") chunks = [refined_text[i:i+50] for i in range(0, len(refined_text), 50)] # 50๊ธ์ ๋จ์ ์ฒญํน # 3. VectorDB ์ปค๋ฅ์ ์๋ฆฝ ๋ฐ ์ปฌ๋ ์ ์ด๊ธฐํ qdrant_client = QdrantClient(url=QDRANT_URL) # ์ปฌ๋ ์ ์ด ์์ผ๋ฉด ๊ณ ์ฐจ์(์: 384์ฐจ์) HNSW ์ธ๋ฑ์ค ๊ธฐ๋ฐ์ผ๋ก ์์ฑ if not qdrant_client.collection_exists(collection_name=COLLECTION_NAME): qdrant_client.create_collection( collection_name=COLLECTION_NAME, vectors_config=VectorParams(size=5, distance=Distance.COSINE), # ์์ ๋ฅผ ์ํด 5์ฐจ์์ผ๋ก ์ค์ ) # 4. ๊ฐ์ ๊ฐ๋ฒผ์ด ์๋ฒ ๋ฉ ๋ฒกํฐ ์์ฑ ๋ฐ Upsert (Idempotency ๋ณด์ฅ) # ์ค๋ฌด์์๋ Ollama๋ SentenceTransformer ๋ชจ๋ธ์ ์์ปค ๋ด๋ถ ํน์ ์ธ๋ถ GPU ๊ฐ์๊ธฐ๋ฅผ ํตํด ์๋ฆฝ points = [] for idx, chunk in enumerate(chunks): dummy_embedding = [0.1 * (idx + 1), 0.2, 0.5, 0.1, 0.9] # 5์ฐจ์ ๊ฐ์ ์๋ฒ ๋ฉ ๋ฒกํฐ points.append(PointStruct( id=idx + 100, # ๊ณ ์ ID ์ง์ ์ผ๋ก ๋ฉฑ๋ฑ์ฑ(Upsert) ํ๋ณด vector=dummy_embedding, payload={"page_content": chunk, "source": "minio://raw/manual_01.txt"} )) qdrant_client.upsert(collection_name=COLLECTION_NAME, points=points) print(f"[Success] Spark ์ ์ฒ๋ฆฌ ์๋ฃ๋ {len(chunks)}๊ฐ์ ์ฒญํฌ๊ฐ Qdrant HNSW ์ธ๋ฑ์ค์ ๋๊ธฐํ๋์์ต๋๋ค.") # ============================================================ # Airflow DAG ์ค์ผ์คํธ๋ ์ด์ ํต์ฌ ์ค์ # ============================================================ default_args = { 'owner': 'ai_platform_eng', 'depends_on_past': False, 'start_date': datetime(2026, 7, 8), 'retries': 2, 'retry_delay': timedelta(minutes=5), } with DAG( dag_id='advanced_bigdata_ai_pipeline_v1', default_args=default_args, description='์์ง->Lake->Spark ์ ์ฒ๋ฆฌ->VectorDB ํ์ดํ๋ผ์ธ ์๋ํ', schedule='@daily', catchup=False, tags=['spark', 'minio', 'qdrant', 'rag'] ) as dag: task_ingest = PythonOperator( task_id='ingest_to_data_lake', python_callable=fn_ingest_to_lake, ) task_spark_and_vector = PythonOperator( task_id='spark_transform_and_vector_upsert', python_callable=fn_spark_processing_and_vector_upsert, ) # ํ์ดํ๋ผ์ธ ์ํ ๊ด๊ณ ๋ฐ์ธ๋ฉ task_ingest >> task_spark_and_vector
- ์ค๋ฌด ๋ชจ๋ํฐ๋ง ๋ฐ ์ํคํ ์ฒ์ ํต์ฌ ํ
- Spark ๋ฉ๋ชจ๋ฆฌ ํ๋ (
OOM๋ฐฉ์ง):
- ์ ์ฒ๋ฆฌ ์ค
SparkDriver๋ ธ๋๋ก ๋๋ฌด ๋ง์ ๋์ฉ๋ ํ ์คํธ ์ง๊ณ ๋ฐ์ดํฐ๋ฅผ ํ ๋ฒ์ ๊ฐ์ ธ์ค๋collect()์ฐ์ฐ์ ์ ๋ ํผํด์ผ ํจ- ๋์ ๊ฐ๊ณต๋ ๋ฐ์ดํฐ๋ฅผ ๊ณง๋ฐ๋ก MinIO/S3์ ๋ถ์ฐ ํ์ผ(Parquet, ๋ฐ์ดํฐ ๋ ์ดํฌ ํฌ๋งท) ํํ๋ก writeํ๋๋ก ํ์ดํ๋ผ์ธ์ ์ค๊ณํด์ผ ํจ
- VectorDB ๋ฐฑํ(Backfill) ๊ณผ๋ถํ ํต์ :
- ๊ณผ๊ฑฐ ๋์ฉ๋ ๋ฐ์ดํฐ๋ฅผ ํ ๋ฒ์ ์ฌ์ฒ๋ฆฌํ ๋, VectorDB์ ์์ฒ๋ง ๊ฑด์ ๋ฒกํฐ๊ฐ ํ๊บผ๋ฒ์ ์์์ง๋ฉด ์ค์๊ฐ HNSW ๊ทธ๋ํ ์ฐ์ฐ ๋๋ฌธ์ DB CPU๊ฐ ๋ง๋น๋ ์ ์์
- ์ด ๊ฒฝ์ฐ Airflow์
max_active_runs_per_dag=1์ค์ ์ ํตํด ๋ฐฐ์น ์คํ ๋จ์๋ฅผ ๊ฐ์ ์กฐ์จํ์ฌ ์์คํ ์์ ๋ง์ ๊ฐ๋ํด์ผ ํจ
3. ์ข ํฉ ์ค์ต ์์ ์ฝ๋ 2
3.1 ์ํคํ ์ฒ ์ค๋ฌด ๋น์ฆ๋์ค ์๋๋ฆฌ์ค
โ์ค๋งํธํฉํ ๋ฆฌ ๊ฐ๋ณ ๋ฉด์ ์ ๋๊ณ(Variable Area Flowmeter) ์ผ์ ์คํธ๋ฆฌ๋ฐ ๋ฐ์ดํฐ์ ํ์ด๋ธ๋ฆฌ๋ ๋ ์ดํฌํ์ฐ์ค ๊ตฌ์ถ ๋ฐ RAG ์ง์ ๋ฒ ์ด์ค ๋๊ธฐํโ
- ์ค์๊ฐ ์์ง (Kafka):
- ๊ณต์ฅ ์ผ์ ๋ฐ ์ค๋น ์ ์ด๊ธฐ์์ ๋ฐ์ํ๋ ์ค์๊ฐ ๋น์ ํ ๋ก๊ทธ ๋ฐ ๋งค๋ด์ผ ํ ์คํธ๊ฐ
- Kafka ํ ํฝ์ผ๋ก ์ธ๋ฑ์ฑ
- ๋ฐ์ดํฐ ๋ ์ดํฌ ๋ณด์กด (MinIO):
- ์์ง๋ ์์(Raw) ๋ก๊ทธ๋ ๋ฐ์ดํฐ ์ ์ค ๋ฐฉ์ง ๋ฐ ๋ฉฑ๋ฑ์ฑ ํ๋ณด๋ฅผ ์ํด
- MinIO ์ค๋ธ์ ํธ ์คํ ๋ฆฌ์ง์ ๋ ์ง๋ณ ํํฐ์
์์ญ(
raw/{ { ds_nodash }}/)์ ์๊ตฌ ๋ณด์กด
- ๋ถ์ฐ ์ ์ฒ๋ฆฌ ๋ฐ ์๋ฒ ๋ฉ ๊ฐ๊ณต (PySpark):
- Spark ์ธ์ ์ ์ปจํ ์ด๋ ๋ด๋ถ์์ ๋์ ๊ตฌ๋ํ์ฌ,
- ๋น์ ํ ํ ์คํธ์ ๋ ธ์ด์ฆ๋ฅผ ์ ๊ฑฐํ๊ณ
- ์๋ฏธ๋ก ์ ๋ฌธ๋งฅ ๋ณด์ ์ ์ํ ์ฒญํน(Chunking) ๋ถ์ฐ ์ฐ์ฐ์ ์ํ
- ์ง์ ๋ฒ ์ด์ค ๋๊ธฐํ (Qdrant):
- ๊ฐ๊ณต๋ ๊ณ ์ฐจ์ ๋ฒกํฐ ๋ฐ์ดํฐ๋ฅผ ๊ณ ์ ๋ฐ์ง ๊ฒ์์ ์ํด
- Qdrant ๋ฒกํฐ ์คํ ์ด์ HNSW ๊ทธ๋ํ ์ธ๋ฑ์ค์ ์ค๋ณต ์์ด Upsertํ์ฌ
- ์ค์๊ฐ RAG ๊ฒ์ ์์ง์ ์ต์ ํํจ
3.2 ์ปจํ ์ด๋ ํ๊ฒฝ ์์ฑ
- โKafka โ MinIO โ Apache Spark โ Qdrantโ ํ์คํ ๋ฐ์ดํฐ ํ๋ซํผ ์ธํ๋ผ ์คํ
- ์ปดํฌ๋ํธ ๊ฐ ๊ฒฉ๋ฆฌ ์ฅ๋ฒฝ ์์ด ์ ๊ธฐ์ ์ผ๋ก ์ฐ๋๋๋๋ก ๋์ผํ ๋ธ๋ฆฟ์ง ๋คํธ์ํฌ(
bigdata-network)๋ก ์ฐ๊ฒฐ
Dockerfile.airflow- Airflow ๊ณต์ docker-compose.yml์ผ๋ก ์ค์น ์, Java/JVM์ ๊ด๋ จ๋ ๋ถ๋ถ์ด ๋ฒ์ ์ฐจ์ด ๋ฑ ๋ช ๊ฐ์ง ๋ฌธ์ ๋ก ์ธํด ์ ๋๋ก ์ค์น๋์ด ์์ง ์์
- ์ค์ต ์๋๋ฆฌ์ค์ ํ์ํ ๊ฐ ์ปจํ ์ด๋๋ฅผ ์ฐ๋, ์คํํ๋ ค๋ฉด OpenJDK 21 ์ด์ ๋ฒ์ ์ด ์๊ตฌ๋จ
- ํนํ OpenJDK 21 ๋ด๋ถ์ GLIBC 2.38 ์ด์์ ์๊ตฌํจ
- ๋ฐ๋ผ์ Airflow 3.3.0์ ๊ธฐ๋ฐ์ผ๋ก ์ปค์คํ
๋น๋๋ฅผ ์ํํ์ฌ์ผ ํจ
- ํฅํ์ ์ฌ์ฉ๋ ๊ฐ ๋ชจ๋ ๋ฐ ๋ผ์ด๋ธ๋ฌ๋ฆฌ ํจํค์ง๋ค๋ ์ปค์คํ ๋น๋์ ๋ฏธ๋ฆฌ ๋ฃ์ด๋์ผ๋ฉด ์ปจํ ์ด๋์ Up/Down์ด ๋นจ๋ผ์ง
# 1. ์ํ์น ์์ดํ๋ก์ฐ ๊ณต์ ๋ฒ ์ด์ค ์ด๋ฏธ์ง ์ง์ FROM apache/airflow:3.3.0 # 2. ์์คํ ํจํค์ง ์ค์น๋ฅผ ์ํด root ๊ถํ์ผ๋ก ์ ์ ์ค์์นญ USER root # 3. ์ปจํ ์ด๋ ๋ด๋ถ OS ํ๊ฒฝ์ 100% ๋ง๋ ์์ OpenJDK 21 ์ค์น (Glibc ์ถฉ๋ ์์ฒ ์ฐจ๋จ) RUN curl -Lf https://github.com/adoptium/temurin21-binaries/releases/download/jdk-21.0.2%2B13/OpenJDK21U-jdk_x64_linux_hotspot_21.0.2_13.tar.gz -o /tmp/openjdk.tar.gz && \ mkdir -p /usr/lib/jvm/java-21-openjdk-amd64 && \ tar -xzf /tmp/openjdk.tar.gz -C /usr/lib/jvm/java-21-openjdk-amd64 --strip-components=1 && \ rm -rf /tmp/openjdk.tar.gz # 4. ์์ดํ๋ก์ฐ ์คํ ๊ถํ(์์ ์)์ผ๋ก ๋ค์ ๋ณต๊ท USER airflow # 5. ๊ธฐ์กด์ _PIP_ADDITIONAL_REQUIREMENTS๋ก ์ค์๊ฐ ๋ค์ด๋ก๋๋ฐ๋ ๋น ๋ฐ์ดํฐ ํจํค์ง๋ค์ ์ด๋ฏธ์ง์ ๋ฏธ๋ฆฌ ๋น๋ RUN pip install --no-cache-dir pyspark==4.1.2 minio qdrant-client kafka-python
- Airflow ๊ณต์ docker-compose.yml์ผ๋ก ์ค์น ์, Java/JVM์ ๊ด๋ จ๋ ๋ถ๋ถ์ด ๋ฒ์ ์ฐจ์ด ๋ฑ ๋ช ๊ฐ์ง ๋ฌธ์ ๋ก ์ธํด ์ ๋๋ก ์ค์น๋์ด ์์ง ์์
docker-compose.yml- Airflow ๊ณตํต ๋ถ๋ถ ์์
- ์์ ์ง์ : build ๋ถ๋ถ, environment์ JAVA_HOME ๋ถ๋ถ
x-airflow-common: &airflow-common image: ${AIRFLOW_IMAGE_NAME:-apache/airflow:3.3.0} build: context: . dockerfile: Dockerfile.airflow env_file: - ${ENV_FILE_PATH:-.env} environment: &airflow-common-env AIRFLOW__CORE__EXECUTOR: CeleryExecutor AIRFLOW__CORE__AUTH_MANAGER: airflow.providers.fab.auth_manager.fab_auth_manager.FabAuthManager AIRFLOW__DATABASE__SQL_ALCHEMY_CONN: postgresql+psycopg2://airflow:airflow@postgres/airflow AIRFLOW__CELERY__RESULT_BACKEND: db+postgresql+psycopg2://airflow:airflow@postgres/airflow AIRFLOW__CELERY__BROKER_URL: redis://:@redis:6379/0 AIRFLOW__CORE__FERNET_KEY: ${FERNET_KEY} AIRFLOW__CORE__DAGS_ARE_PAUSED_AT_CREATION: 'true' AIRFLOW__CORE__LOAD_EXAMPLES: 'true' AIRFLOW__CORE__EXECUTION_API_SERVER_URL: 'http://airflow-apiserver:8080/execution/' AIRFLOW__API_AUTH__JWT_SECRET: ${AIRFLOW__API_AUTH__JWT_SECRET:-airflow_jwt_secret} AIRFLOW__API_AUTH__JWT_ISSUER: ${AIRFLOW__API_AUTH__JWT_ISSUER:-airflow} AIRFLOW__SCHEDULER__ENABLE_HEALTH_CHECK: 'true' _PIP_ADDITIONAL_REQUIREMENTS: ${_PIP_ADDITIONAL_REQUIREMENTS:-} AIRFLOW_CONFIG: '/opt/airflow/config/airflow.cfg' JAVA_HOME: "/usr/lib/jvm/java-21-openjdk-amd64"
- Kafka ์ปจํ
์ด๋ ์ถ๊ฐ
- ๊ธฐ์กด์๋ Kafka์ ๋ํ ์ ๊ทผ ๋๋ฉ์ธ์ ์ธ๋ถ, ๋ด๋ถ๋ก ๋๋์ด ์ฒ๋ฆฌํ์์ผ๋,
- ํ์ฌ๋ ์ธ๋ถ ์ ๊ทผ๋ ๊ฒฐ๊ตญ Airflow๊ฐ ์ฒ๋ฆฌํ๋ฏ๋ก ๋ด๋ถ์ ๋์ผํ๊ฒ ๋๋ฉ์ธ์ ์ค์ ํจ
- ๊ทธ๋ฌ๋ ํฅํ ์ ์ ์๋น์ค๋ก์ ํ์ฅ ์, ๋ณด์ ๋ฑ์ ๊ณ ๋ คํ์ฌ ํฌํธ์ ๋ณํ์ ๊ทธ๋๋ก ์ ์งํจ
- (ํ์ฌ ์์ ์์๋ ๊ตฌ๋ถํ๋ ์๋ฏธ๊ฐ ์์)
kafka-1: image: apache/kafka:4.3.1 container_name: kafka-1 ports: - "9092:9092" environment: KAFKA_NODE_ID: 1 KAFKA_PROCESS_ROLES: broker,controller KAFKA_LISTENERS: INTERNAL://0.0.0.0:19092, EXTERNAL://0.0.0.0:9092, CONTROLLER://0.0.0.0:9093 KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka-1:19092,EXTERNAL://kafka-1:9092 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka-1:9093,2@kafka-2:9093,3@kafka-3:9093 KAFKA_LOG_DIRS: /var/lib/kafka/data volumes: - kafka-1-data-volume:/var/lib/kafka/data kafka-2: image: apache/kafka:4.3.1 container_name: kafka-2 ports: - "9094:9092" environment: KAFKA_NODE_ID: 2 KAFKA_PROCESS_ROLES: broker,controller KAFKA_LISTENERS: INTERNAL://0.0.0.0:19092, EXTERNAL://0.0.0.0:9092, CONTROLLER://0.0.0.0:9093 KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka-2:19092,EXTERNAL://kafka-2:9092 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka-1:9093,2@kafka-2:9093,3@kafka-3:9093 KAFKA_LOG_DIRS: /var/lib/kafka/data volumes: - kafka-2-data-volume:/var/lib/kafka/data kafka-3: image: apache/kafka:4.3.1 container_name: kafka-3 ports: - "9095:9092" environment: KAFKA_NODE_ID: 3 KAFKA_PROCESS_ROLES: broker,controller KAFKA_LISTENERS: INTERNAL://0.0.0.0:19092, EXTERNAL://0.0.0.0:9092, CONTROLLER://0.0.0.0:9093 KAFKA_ADVERTISED_LISTENERS: INTERNAL://kafka-3:19092,EXTERNAL://kafka-3:9092 KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: INTERNAL:PLAINTEXT,EXTERNAL:PLAINTEXT,CONTROLLER:PLAINTEXT KAFKA_INTER_BROKER_LISTENER_NAME: INTERNAL KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka-1:9093,2@kafka-2:9093,3@kafka-3:9093 KAFKA_LOG_DIRS: /var/lib/kafka/data volumes: - kafka-3-data-volume:/var/lib/kafka/data
-
MinIO ์ปจํ ์ด๋ ์ถ๊ฐ
minio: image: minio/minio:RELEASE.2025-04-22T22-12-26Z container_name: minio ports: - "9000:9000" # ํ์ด์ฌ SDK๊ฐ ์ ์ํ API ํต์ ํฌํธ - "9001:9001" # ์น UI ์ฝ์ ์ด๋๋ฏผ ํฌํธ environment: MINIO_ROOT_USER: minioadmin MINIO_ROOT_PASSWORD: minioadmin volumes: - minio-data-volume:/data command: server /data --console-address ":9001"
-
Qdrant ์ปจํ ์ด๋ ์ถ๊ฐ
qdrant: image: qdrant/qdrant:v1.18.0 container_name: qdrant ports: - "6333:6333" # REST API์ฉ ํฌํธ (์์ดํ๋ก์ฐ ์์ปค ํต์ ์ฑ๋) - "6334:6334" # gRPC์ฉ ํฌํธ volumes: - qdrant-data-volume:/qdrant/storage
-
๊ฐ ์ปจํ ์ด๋์์ ์ฌ์ฉํ๋ ๋ณผ๋ฅจ ์ถ๊ฐ
volumes: kafka-1-data-volume: kafka-2-data-volume: kafka-3-data-volume: postgres-db-volume: minio-data-volume: qdrant-data-volume:
- Airflow ๊ณตํต ๋ถ๋ถ ์์
3.3 ์ฌ์ ํ๊ฒฝ ๊ตฌ์ถ ๋ฐ ํ ํด๋ก์ง ์ค์
- ์ธ๋ถ์์ ๊ตฌ๋ ์ค์ธ ๋ํ ์คํ์์ค ์ธํ๋ผ ์ปดํฌ๋ํธ๋ค์ ์ปจํ ์ด๋ ๋ด๋ถ์ Airflow๊ฐ ์ธ์ํ๊ณ ,
- Python 3.13 ํ๊ฒฝ์์ ๊ด๋ จ ์์กด์ฑ ํฌ๋์ ์์ด ํจํค์ง๋ฅผ ๋ก๋ํ ์ ์๋๋ก
.env๋ฅผ ์ ์
# 1. ๊ธฐ์กด ๊ตฌ๋ฒ์ ์ธํ๋ผ ์ปจํ
์ด๋ ์์ ์ญ์ (๋ณผ๋ฅจ ๋ณด์กด ์ ํ ๊ฐ๋ฅํ๋ ํด๋ฆฐ ์ํคํ
์ฒ๋ฅผ ์ํด ๋ค์ด ์ถ์ฒ)
docker compose down -v
# 2. ๋ก์ปฌ ๊ฐ์ ํ๊ฒฝ ๋ณ์ ์ ๊ฒ (.env์ ๋ช
์๋ ์ต์ ํจํค์ง ํ์ธ)
# _PIP_ADDITIONAL_REQUIREMENTS=minio qdrant-client pyspark kafka-python trino
echo -e "AIRFLOW_UID=$(id -u)" > .env
# 3. ์์ดํ๋ก์ฐ ์ต์ ์ด๋ฏธ์ง(3.3.0) ๊ธฐ๋ฐ ๋ฉํ ์คํค๋ง ๋ง์ด๊ทธ๋ ์ด์
๋ฐ ์ด๊ธฐ ๊ณ์ ์ค์
docker compose up airflow-init
# 4. KRaft ๋ฐ ์ง์ MinIO ๋ฒ์ ์ด ํตํฉ๋ ์ ์ฒด ์ปดํฌ๋ํธ ์ค์๊ฐ ๋ฐฑ๊ทธ๋ผ์ด๋ ๊ตฌ๋
docker compose up -d
# 5. ์ ์ ๊ธฐ๋ ์ํ ๊ฒ์ฆ
docker compose ps
- ํฌํธ ๋งตํ ์ต์ข
์ฃผ์:
- Airflow ์ธํ๋ผ ํฌํ:
http://localhost:8080(๊ณ์ :airflow/airflow) - MinIO Console:
http://localhost:9001(๊ณ์ :minioadmin/minioadmin) - Qdrant Dashboard:
http://localhost:6333/dashboard
- Airflow ์ธํ๋ผ ํฌํ:
3.5 DAG ๊ตฌ์ฑ
- ํ์ดํ๋ผ์ธ ์์กด์ฑ ์ํคํ
์ฒ ๋ฐ ํ๋ฆ๋
- ๊ฐ ์คํ์์ค ์ปดํฌ๋ํธ์ ๊ฒฐํจ ๊ฒฉ๋ฆฌ(Fault Isolation)๋ฅผ ์๊ฐํํ ์์ ์ ๋ ฌ ์์กด์ฑ ๊ตฌ์กฐ
-
๋ฐ์ดํฐ ์ ์ด ํ๋ฆ ๊ตฌ์กฐ (Data Control Flow)
- ๋จ๊ณ๋ณ ๋งค์ปค๋์ฆ
task_kafka_kr_ingest(์์ง ๊ณ์ธต):- KRaft ๋จ๋ ๋ ธ๋๋ก ๊ธฐ๋ ์ค์ธ Kafka ํ ํฝ์ผ๋ก
- ๊ฐ๋ณ ๋ฉด์ ์ ๋๊ณ ์ฅ์ ๋ก๊ทธ ์ด๋ฒคํธ๋ฅผ ๋ฐํํ๊ณ ์์ง ๊ฒ์ฆ
- ๋ณ๋ ฌ ์ฒ๋ฆฌ ๊ณ์ธต (Parallel Operations):
- ์์ง ์๋ฃ ์ ํธ๋ฅผ ๋ฐ์ผ๋ฉด,
- ๋ฐ์ดํฐ ๋ ์ดํฌ ๋ฐฑ์ (MinIO), ๋ถ์ฐ ๋ฉ๋ชจ๋ฆฌ ์ ์ฒ๋ฆฌ(Spark)๊ฐ
- ํธ์คํธ ์์์ ํจ์จ์ ์ผ๋ก ๋๋์ด ์ฐ๋ฉฐ ๋์์ ๋ณ๋ ฌ๋ก ๊ธฐ๋
task_vector_upsert_qdrant(์ต์ข ์ ์ฌ ๊ณ์ธต):- ์์ ์ธ ๊ฐ๋์ ๋ฐ์ดํฐ ์์ง๋์ด๋ง ํ์ดํ๋ผ์ธ์ด ๋ชจ๋ ๋ฌด์ฌํ ์ฑ๊ณต(
Success) ์ํ๋ก ๋งํน๋์ด์ผ๋ง - ์ต์ข RAG ๊ฒ์์ ์ํ Qdrant ๋ฒกํฐ์คํ ์ด ์ ์ํธ๋ฅผ ๋จํ
- ์์ ์ธ ๊ฐ๋์ ๋ฐ์ดํฐ ์์ง๋์ด๋ง ํ์ดํ๋ผ์ธ์ด ๋ชจ๋ ๋ฌด์ฌํ ์ฑ๊ณต(
- ๋จ๊ณ๋ณ ๋งค์ปค๋์ฆ
-
DAG ์์ค ์ฝ๋
- ์ฝ๋ ์ค์ ๊ท์ฝ
- PEP 8 ํ์ค(with ๋ฌธ ํ์ 4์นธ ๋ค์ฌ์ฐ๊ธฐ),
- ์ต์ Airflow ๋งค๊ฐ๋ณ์(
schedule='@daily'), - KRaft ๋ฉ์์ง ํ ํต์ ์ธํฐํ์ด์ค
-
~/workspace/airflow/dags/advanced_bigdata_ai_pipeline_v5.py๊ฒฝ๋ก๋ก ์์ฑ#//file: "dags/advanced_bigdata_ai_pipeline_v5.py" import io import json from datetime import datetime, timedelta from airflow import DAG from airflow.providers.standard.operators.python import PythonOperator from minio import Minio from kafka import KafkaProducer from pyspark.sql import SparkSession from qdrant_client import QdrantClient from qdrant_client.models import Distance, VectorParams, PointStruct KAFKA_BROKERS = [ "kafka-1:9092", "kafka-2:9092", "kafka-3:9092" ] MINIO_URL = "minio:9000" MINIO_ACCESS = "minioadmin" MINIO_SECRET = "minioadmin" RAW_LAKE_BUCKET = "factory-raw-stream" COLLECTION_NAME = "factory_hybrid_knowledge_base" QDRANT_HOST = "qdrant" QDRANT_PORT = "6333" def fn_kafka_kr_stream_ingest(**context): sample_log = { "equipment": "Variable Area Flowmeter (๊ฐ๋ณ ๋ฉด์ ์ ๋๊ณ)", "status": "Warning", "message": "์ ๋ ๋ณ๋์ ๋ฐ๋ฅธ ํ๋ธ ๋ด ํ๋กํธ(Float) ์ง๋ ๋ฐ์. ๊ฐ๋ณ ๋ฉด์ ์ธก์ ์ ๋ฐ๋ ์ ํ ์ฐ๋ ค. ๊ธฐ๋ฐ ํจํน ์ค๋ง(O-ring) ๊ต์ฒด ์๋ง." } try: # ๋ฉํฐ ๋ ธ๋ ์ฟผ๋ผ์ ๋ฉ์์ง๋ฅผ ์์ ํ๊ฒ ๋ถ์ฐ ๋ฐํํ๊ธฐ ์ํ ๊ณ ๊ฐ์ฉ์ฑ ํ๋ก๋์ ๋น๋ producer = KafkaProducer( bootstrap_servers=KAFKA_BROKERS, # 3๊ฐ ๋ธ๋ก์ปค ํ ์ฃผ์ value_serializer=lambda v: json.dumps(v).encode("utf-8"), max_block_ms=5000, # ์ฟผ๋ผ ํฉ์ ๋๊ธฐ ๋ง์ง ์ํญ ํ์ฅ acks='all', # ๋ฆฌ๋์ ํ๋ก์ ๋ธ๋ก์ปค ๋ชจ๋ ์ ์ฌ ์ฑ๊ณต ํฉํธ ํ์ธ ๊ท์น retries=3 # ์ผ์์ ๋ ธ๋ ์๋ฒ ๊ต์ฒด ์ ์์ฒด ์ฌ์๋ ์์ ์ฅ์น ) # ํ ํฝ ๋ฐํ ์์ ์ ํํฐ์ ๋ถ์ฐ ์ ์ฌ๊ฐ ๊ฐ๋ฅํ๋๋ก ์ ๋๊ณ ๋ฉ์์ง ํฌํ producer.send('factory-sensor-topic', value=sample_log) producer.flush() print("\n" + "="*80) print(f"[Apache Kafka 4.3.1] 3๊ฐ ๋ ธ๋ ๋ฉํฐ ํด๋ฌ์คํฐ ์ธํ๋ผ๋ง({KAFKA_BROKERS})์ ์คํธ๋ฆฌ๋ฐ ์ด๋ฒคํธ๋ฅผ ์ ์ค ์์ด ์ฑ๊ณต์ ์ผ๋ก ๋ฐํํ์ต๋๋ค.") print("\n" + "="*80) except Exception as e: print(f"[๋คํธ์ํฌ ์ฐํ ๊ฐ์ด๋] ๋ฉํฐ ์นดํ์นด ์ฟผ๋ผ ์ฐ๊ฒฐ ์ฐจ๋จ์ผ๋ก ๋ด์ฅ ์ปจํ ์คํธ ์คํธ๋ฆผ ๋ฉ๋ชจ๋ฆฌ๋ก ๋์ฒด ๊ตฌ๋ํฉ๋๋ค. ์ฌ์ : {e}") context["ti"].xcom_push(key="raw_stream_data", value=sample_log) def fn_minio_raw_backup(**context): ds_nodash = context['ds_nodash'] raw_data = context["ti"].xcom_pull(task_ids="task_kafka_kr_ingest", key='raw_stream_data') client = Minio(MINIO_URL, access_key=MINIO_ACCESS, secret_key=MINIO_SECRET, secure=False) if not client.bucket_exists(RAW_LAKE_BUCKET): client.make_bucket(RAW_LAKE_BUCKET) object_path = f"raw/{ds_nodash}/stream_log.json" client.put_object( bucket_name=RAW_LAKE_BUCKET, object_name=object_path, data=io.BytesIO(json.dumps(raw_data, ensure_ascii=False).encode("utf-8")), length=len(json.dumps(raw_data, ensure_ascii=False).encode("utf-8")), content_type="application/json" ) print("\n" + "="*80) print(f"[MinIO ๋ฐฑ์ ์ฑ๊ณต] ์ง์ ๋ฒ์ ์คํ ๋ฆฌ์ง ์ ์ฌ ์๋ฃ ๊ฒฝ๋ก: {object_path}") print("\n" + "="*80) def fn_spark_transform_processing(**context): raw_data = context["ti"].xcom_pull(task_ids="task_kafka_kr_ingest", key="raw_stream_data") spark = SparkSession.builder \ .appName("KRaftEnvironmentSparkProcessor") \ .master("local[*]") \ .getOrCreate() target_corpus = f"[๊ณ์ธก๊ธฐ ์ค์๊ฐ ์ฅ์ ๊ฐ์ด๋] ๋์ ์ค๋น: {raw_data['equipment']} | ํ์ฅ ์ํ: {raw_data['status']} | ์กฐ์น ์ง์นจ: {raw_data['message']}" df = spark.createDataFrame([(1, target_corpus)], ["id", "text"]) refined_text = df.select("text").collect()[0][0] spark.stop() print("\n" + "="*80) print(f"[PySpark ๋ถ์ฐ ๊ฐ๊ณต ์์ง ์ฑ๊ณต] ๊ฐ๋ณ ๋ฉด์ ์ ๋๊ณ ์คํธ๋ฆฌ๋ฐ ํ ์คํธ ์ ์ ์์!") print(f" โ ์ ์ ๋ ์ฒญํฌ: {refined_text}") print("="*80 + "\n") context["ti"].xcom_push(key="refined_chunk", value=refined_text) def fn_vector_upsert_qdrant(**context): ds_nodash = context['ds_nodash'] processed_chunk = context["ti"].xcom_pull(task_ids="task_spark_transform", key="refined_chunk") qdrant_client = QdrantClient(host=QDRANT_HOST, port=QDRANT_PORT) if not qdrant_client.collection_exists(collection_name=COLLECTION_NAME): qdrant_client.create_collection( collection_name=COLLECTION_NAME, vectors_config=VectorParams(size=3, distance=Distance.COSINE) ) point_id = int(f"{ds_nodash}99") qdrant_client.upsert( collection_name=COLLECTION_NAME, points=[ PointStruct( id=point_id, vector=[0.18, 0.65, 0.49], payload={ "page_content": processed_chunk, "log_date": ds_nodash, "lineage": "kafka ๐กช minio ๐กช pyspark ๐กช qdrant" } ) ] ) print("\n" + "="*80) print(f"[Qdrant] ์ต์ RAG ์ง์ ๋ฒ ์ด์ค ๋์ ๋ฐ์ดํฐ Upsert ์์ ์๋ฃ (Point ID: {point_id})") print("\n" + "="*80) default_args = { "owner": "seokhwan", "depends_on_past": False, "start_date": datetime(2026, 7, 11), "retries": 1, "retry_delay": timedelta(minutes=3) } with DAG( dag_id="advanced_bigdata_ai_pipeline_v5", default_args=default_args, description="์ต์ ์ํ์น ์นดํ์นด 4.3.1 ๋ฉํฐ ๋ธ๋ก์ปค ๋ฐ ์ง์ MinIO ๊ธฐ๋ฐ ์ํฐํ๋ผ์ด์ฆ ์๋ํ ํ์ดํ๋ผ์ธ", schedule="@daily", catchup=False, tags=["kafka", "minio", "spark", "qdrant"] ) as dag: task_kafka_kr_ingest = PythonOperator( task_id="task_kafka_kr_ingest", python_callable=fn_kafka_kr_stream_ingest ) task_minio_raw_backup = PythonOperator( task_id="task_minio_raw_backup", python_callable=fn_minio_raw_backup ) task_spark_transform = PythonOperator( task_id="task_spark_transform", python_callable=fn_spark_transform_processing ) task_vector_upsert_qdrant = PythonOperator( task_id="task_vector_upsert_qdrant", python_callable=fn_vector_upsert_qdrant ) task_kafka_kr_ingest >> [task_minio_raw_backup, task_spark_transform] >> task_vector_upsert_qdrant
- ์ฝ๋ ์ค์ ๊ท์ฝ
3.5 ๊ฒฐ๊ณผ ํ์ธ
-
http://localhost:8080์ ์ ์ ํ DAG ๋ฉ๋ด์์enterprise_bigdata_ai_stream_pipeline_v5DAG๊ฐ ๋ํ๋๋์ง ๊ฒ์์ฐฝ์์ ์กฐํ -
๊ฒฐ๊ณผ ํ์ธ
- ํ์ธ ์ ์ํฉ
- ๊ธฐ์กด์ ์คํจํ๋ ์ค์ ์ดํ๋ก ์ฑ๊ณตํ ์ค์ ์ ํ์ธํ ์ ์์
- ์ด๋ฐ์ ์ฑ๊ณต ์ดํ, ์ค์ ์ด ์์ ํ๋ ๋ค์๋ ์ฒ๋ฆฌ ์๊ฐ๋ ํฌ๊ฒ ๊ฐ์ํ์์
- Trigger ์คํ ๊ฒฐ๊ณผ
- ์ ๊ณผ์ ์ด ๋ฌด๋ฆฌ์์ด ์ฑ๊ณตํจ
- ์ฒ๋ฆฌ ์๊ฐ: ์์๋ถํฐ ๋๊น์ง ์ฝ 8์ด์ ์๋ฃ๋์์
- ์ค์ ์์ ์ ๋ฐ๋ฅธ ์์ ํ ํ ์์ ์๊ฐ์ด ๊ธ๊ฐํ์์ ํ์ธ
- task_kafka_kr_ingest์ ๋ก๊ทธ
- 3๊ฐ์ Kafka ๋ณผ๋ฅจ์ ์ ๋๋ก ๋ณต์ , ์ ์ฅ๋์์
- task_spark_transform์ ๋ก๊ทธ
- ๊ธฐ์กด์ ์คํจํ๋ ์ค์ ์ดํ๋ก ์ฑ๊ณตํ ์ค์ ์ ํ์ธํ ์ ์์
- ์ค๊ฐ์ Error ํ์๊ฐ ๋ณด์ด์ง๋ง ๋ด์ฉ์ ์ฝ์ด๋ณด๋ฉด Error๊ฐ ์๋์ ์ ์ ์์
- ERROR ํ์๊ฐ ์ถ๋ ฅ๋ ๊ฒ์ ์คํํฌ์ ํ์ค ์๋ฐ ์๋ฌ ์ถ๋ ฅ(Stderr) ์คํธ๋ฆผ์ด ์ ์
๋์ด ๋ฐ์ํ ์์ดํ๋ก์ฐ ๊ณ ์ ์ ๋ก๊น
ํน์ง
- ์คํํฌ ์์ง์ ๊ตฌ๋๋ ๋ ์์ง์ ์์คํ ๊ฒฝ๊ณ ๋ ํ๊ฒฝ ์ค์ ์๋ด๋ฅผ Stdout์ด ์๋ Stderr ์คํธ๋ฆผ์ผ๋ก ๋ด๋ณด๋ด๋๋ก ์ฝ๋๊ฐ ์์ฑ๋์ด ์์
- ์์ดํ๋ก์ฐ๊ฐ ์ถ๋ ฅ ์คํธ๋ฆผ์ด Stderr๋ก ์ ๋ฌ๋ ๊ฒ๋ง ๋ณด๊ณ ๊ธฐ๊ณ์ ์ผ๋ก ERROR๋ฅผ ๋ถ์ฌ๋ฒ๋ฆฐ ๊ฒ์
- ์ง์ง ๋ด๋ถ ์์ธ(Exception)๊ฐ ๋ฐ์ํ๋ค๋ฉด ์คํํฌ ํน์ ์ ๊ฑฐ๋ํ StackTrace ์๋ฐ ์๋ฌ ๋ฌธ๋จ์ด ์ถ๋ ฅ๋์ด์ผ ํจ
- ERROR ํ์๊ฐ ์ถ๋ ฅ๋ ๊ฒ์ ์คํํฌ์ ํ์ค ์๋ฐ ์๋ฌ ์ถ๋ ฅ(Stderr) ์คํธ๋ฆผ์ด ์ ์
๋์ด ๋ฐ์ํ ์์ดํ๋ก์ฐ ๊ณ ์ ์ ๋ก๊น
ํน์ง
- task_minio_raw_backup ๋ก๊ทธ
- ๋ฐ์ดํฐ์ ์ ์ฌ๊ฐ ์ ์์ ์ผ๋ก ์๋ฃ๋์์
- task_vector_upsert_qdrant์ ๋ก๊ทธ
- ๋ฐ์ดํฐ์ Upsert๊ฐ ์ ์์ ์ผ๋ก ์๋ฃ๋์์
- ํ์ธ ์ ์ํฉ