1. ์˜ค์ผ€์ŠคํŠธ๋ ˆ์ด์…˜ ์•„ํ‚คํ…์ฒ˜

  • โ€œ์ˆ˜์ง‘ ๐Ÿกฒ Data Lake ๐Ÿกฒ Apache Spark ๐Ÿกฒ VectorDBโ€๋กœ ์ด์–ด์ง€๋Š” ํŒŒ์ดํ”„๋ผ์ธ:
    • ํ˜„๋Œ€์ ์ธ Enterprise AI ๋ฐ ๋Œ€๊ทœ๋ชจ ํ•˜์ด๋ธŒ๋ฆฌ๋“œ RAG(Retrieval-Augmented Generation) ์‹œ์Šคํ…œ์˜ ํ‘œ์ค€ ๋ฐฑ์—”๋“œ ์•„ํ‚คํ…์ฒ˜
    • ์‹ค๋ฌด ํ™˜๊ฒฝ์—์„œ๋Š” ํ…Œ๋ผ๋ฐ”์ดํŠธ๊ธ‰ ๋Œ€์šฉ๋Ÿ‰ ๋น„์ •ํ˜• ๋ฐ์ดํ„ฐ(๋ฌธ์„œ, ๋กœ๊ทธ)๊ฐ€ ์œ ์ž…๋˜๋ฏ€๋กœ, ๋‹จ์ผ Airflow ์›Œ์ปค์—์„œ ๋ฐ์ดํ„ฐ๋ฅผ ๊ฐ€๊ณตํ•˜๋Š” ๊ฒƒ์€ ๋ถˆ๊ฐ€๋Šฅ
    • ๋”ฐ๋ผ์„œ Airflow๋Š” ์ œ์–ด(Orchestration)๋งŒ ๋‹ด๋‹นํ•˜๊ณ , ์‹ค์ œ ์ค‘๋Ÿ‰ ์—ฐ์‚ฐ์€ ๋ถ„์‚ฐ ์ปดํ“จํŒ… ์—”์ง„(Spark)์— ์œ„์ž„ํ•˜๋Š” ๊ตฌ์กฐ๋ฅผ ์ทจํ•จ

1.1 ์ •์˜ ๋ฐ ๊ฐœ๋…

  • ์˜ค์ผ€์ŠคํŠธ๋ ˆ์ด์…˜ ์•„ํ‚คํ…์ฒ˜(Orchestration Architecture)
    • ๋ฐ์ดํ„ฐ ์—”์ง€๋‹ˆ์–ด๋ง ๋ฐ AI ์ธํ”„๋ผ์—์„œ ๋ถ„์‚ฐ๋œ ์‹œ์Šคํ…œ, ์„œ๋น„์Šค, ๋ณต์žกํ•œ ๋ฐ์ดํ„ฐ ํŒŒ์ดํ”„๋ผ์ธ์˜
      • ์‹คํ–‰ ์ˆœ์„œ, ์˜์กด์„ฑ, ์ž์› ํ• ๋‹น, ์˜ˆ์™ธ ์ฒ˜๋ฆฌ ๋“ฑ์„ ์ค‘์•™์—์„œ ์ „๋ฐ˜์ ์œผ๋กœ ์กฐ์œจ(Orchestrate)ํ•˜๊ณ  ํ†ต์ œํ•˜๋Š” ํ•ต์‹ฌ ํ”„๋ ˆ์ž„์›Œํฌ๋ฅผ ์˜๋ฏธ
    • ๋‹จ์ˆœํžˆ โ€œ์Šค์ผ€์ค„๋Ÿฌ์— ๋งž์ถฐ ๋ฐฐ์น˜ ์Šคํฌ๋ฆฝํŠธ๋ฅผ ์‹คํ–‰ํ•˜๋Š” ๊ฒƒโ€์„ ๋„˜์–ด,
      • ๊ฑฐ๋Œ€ํ•œ ๋ถ„์‚ฐ ์ปดํ“จํŒ… ํ™˜๊ฒฝ์„ ์•ˆ์ „ํ•˜๊ณ  ๋ฉฑ๋“ฑ์„ฑ ์žˆ๊ฒŒ ๊ด€๋ฆฌํ•˜๊ธฐ ์œ„ํ•œ ๊ณ ๋„์˜ ์‹œ์Šคํ…œ ์•„ํ‚คํ…์ฒ˜
    • ์ˆ˜๋งŽ์€ ๋ถ„์‚ฐ ์ธํ”„๋ผ์™€ ๋ฐ์ดํ„ฐ ์†Œ์Šค๋“ค ์‚ฌ์ด์—์„œ ๋‹ค์Œ์— ๋Œ€ํ•œ ํ•ด๋‹ต์„ ์ œ์‹œํ•˜๋Š” ๊ณ ๋„์˜ ์‹œ์Šคํ…œ ์„ค๊ณ„ ๊ธฐ์ˆ 
      • ์–ด๋–ป๊ฒŒ ๊ฒฐํ•จ์„ ๊ฒฉ๋ฆฌํ•˜๊ณ ,
      • ์–ด๋–ป๊ฒŒ ์ž์›์„ ์ตœ์ ํ™”ํ•˜๋ฉฐ,
      • ์–ด๋–ป๊ฒŒ ๋ฐ์ดํ„ฐ ํ๋ฆ„์„ ์•ˆ์ „ํ•˜๊ฒŒ ๋ณด์žฅํ•  ๊ฒƒ์ธ๊ฐ€
  • ์•„ํ‚คํ…์ฒ˜ ์œ ํ˜•: ํŒจํ„ด ๊ธฐ๋ฐ˜ ๋ถ„๋ฅ˜
    • ์˜ค์ผ€์ŠคํŠธ๋ ˆ์ด์…˜์€ ๋ณต์žกํ•œ ํ•˜๋ถ€ ์„œ๋น„์Šค๋ฅผ ์กฐ์œจํ•˜๋Š” ๋ฐฉ์‹์— ๋”ฐ๋ผ ํฌ๊ฒŒ ๋‘ ๊ฐ€์ง€ ์„ค๊ณ„ ํŒจํ„ด์œผ๋กœ ๋‚˜๋‰จ

    • ์˜ค์ผ€์ŠคํŠธ๋ ˆ์ด์…˜ ํŒจํ„ด (Orchestration Pattern) - ์ค‘์•™ ์ง‘์ค‘ํ˜•
      • ์ค‘์•™์— ๊ฐ•๋ ฅํ•œ โ€˜์ง€ํœ˜์ž(Orchestrator)โ€™ ์—ญํ• ์„ ํ•˜๋Š” ์—”์ง„(์˜ˆ: Apache Airflow, Temporal, Prefect)์„ ๋‘๊ณ ,
      • ์ด ์—”์ง„์ด ๋ชจ๋“  ์„œ๋น„์Šค์™€ ํƒœ์Šคํฌ์˜ ์ƒํƒœ๋ฅผ ์ง์ ‘ ์ œ์–ดํ•˜๊ณ  ๋ช…๋ น์„ ๋‚ด๋ฆฌ๋Š” ๋ฐฉ์‹

      • ์žฅ์ :
        • ์ „์ฒด ํŒŒ์ดํ”„๋ผ์ธ์˜ ์›Œํฌํ”Œ๋กœ์šฐ์™€ ์ƒํƒœ๋ฅผ ์ค‘์•™(Web UI ๋“ฑ)์—์„œ ํ•œ๋ˆˆ์— ๋ชจ๋‹ˆํ„ฐ๋งํ•˜๊ณ  ๊ฐ€์‹œํ™”ํ•  ์ˆ˜ ์žˆ์Œ
        • ์—๋Ÿฌ ๋ฐœ์ƒ ์‹œ ์ค‘์•™์—์„œ ์žฌ์‹œ๋„๋‚˜ ๊ฒฐํ•จ ๊ฒฉ๋ฆฌ๋ฅผ ์ฆ‰๊ฐ ํ†ต์ œํ•  ์ˆ˜ ์žˆ์Œ
      • ๋‹จ์ :
        • ์ค‘์•™ ์˜ค์ผ€์ŠคํŠธ๋ ˆ์ดํ„ฐ ์—”์ง„์ด ๋งˆ๋น„๋˜๊ฑฐ๋‚˜ ๋ฉ”ํƒ€๋ฐ์ดํ„ฐ DB๊ฐ€ ๋‹ค์šด๋˜๋ฉด
          • ์ „์ฒด ์‹œ์Šคํ…œ ํŒŒ์ดํ”„๋ผ์ธ์ด ๋งˆ๋น„๋˜๋Š” ๋‹จ์ผ ์žฅ์• ์ (SPOF, Single Point of Failure)์ด ๋  ์ˆ˜ ์žˆ์Œ
    • ์ฝ”๋ ˆ์˜ค๊ทธ๋ž˜ํ”ผ ํŒจํ„ด (Choreography Pattern) - ์ด๋ฒคํŠธ ๋ถ„์‚ฐํ˜•
      • ์ค‘์•™์˜ ํ†ต์ œ์ž ์—†์ด, ๊ฐ ์„œ๋น„์Šค๋“ค์ด ๋ฉ”์‹œ์ง€ ๋ธŒ๋กœ์ปค(Kafka, RabbitMQ)๋ฅผ ํ†ตํ•ด ์ด๋ฒคํŠธ(Event)๋ฅผ ๋ฐœํ–‰ํ•˜๊ณ  ์ˆ˜์‹ ํ•˜๋ฉฐ
      • ์ž์œจ์ ์œผ๋กœ ์ถค์ถ”๋“ฏ ์ƒํ˜ธ์ž‘์šฉํ•˜๋Š” ๋ฌด์šฉ(Choreography) ๋ฐฉ์‹

      • ์žฅ์ :
        • ์„œ๋น„์Šค ๊ฐ„์˜ ๊ฒฐํ•ฉ๋„(Coupling)๊ฐ€ ๊ทน๋„๋กœ ๋‚ฎ์œผ๋ฉฐ,
        • ํŠน์ • ์„œ๋น„์Šค๊ฐ€ ์ฃฝ์–ด๋„ ๋‹ค๋ฅธ ์„œ๋น„์Šค๋Š” ์ด๋ฒคํŠธ๋ฅผ ๊ณ„์† ์ฒ˜๋ฆฌํ•  ์ˆ˜ ์žˆ์–ด ํ™•์žฅ์„ฑ๊ณผ ๊ฐ€์šฉ์„ฑ์ด ๋›ฐ์–ด๋‚จ
      • ๋‹จ์ :
        • ํŒŒ์ดํ”„๋ผ์ธ์˜ ์ „์ฒด ๋ฐ์ดํ„ฐ ํ๋ฆ„์„ ํ•œ๋ˆˆ์— ํŒŒ์•…ํ•˜๊ธฐ ์–ด๋ ต๊ณ ,
        • ํŠน์ • ๊ตฌ๊ฐ„์—์„œ ๋ฐ์ดํ„ฐ ์ •ํ•ฉ์„ฑ์ด ๊นจ์กŒ์„ ๋•Œ
          • ์—ญ์ถ”์ (๋””๋ฒ„๊น…) ๋ฐ ๋ถ„์‚ฐ ํŠธ๋žœ์žญ์…˜ ๋กค๋ฐฑ(Saga ํŒจํ„ด ๊ตฌํ˜„ ๋“ฑ)์˜ ๋‚œ์ด๋„๊ฐ€ ๋น„์•ฝ์ ์œผ๋กœ ์ƒ์Šน
    • AI ๋ฐ ๋ฐ์ดํ„ฐ ํŒŒ์ดํ”„๋ผ์ธ์˜ ์„ ํƒ:
      • ๋ฐ์ดํ„ฐ ์ˆ˜์ง‘ ๐Ÿกฒ ์ „์ฒ˜๋ฆฌ ๐Ÿกฒ ๋ชจ๋ธ ํ•™์Šต์œผ๋กœ ์ด์–ด์ง€๋Š” ์—„๊ฒฉํ•œ ์„ ํ›„ ๊ด€๊ณ„์™€ ์ธ๊ณผ๊ด€๊ณ„๊ฐ€ ์ค‘์š”ํ•œ
      • AI/๋ฐ์ดํ„ฐ ์—”์ง€๋‹ˆ์–ด๋ง ์˜์—ญ์—์„œ๋Š” ๊ฐ€์‹œ์„ฑ๊ณผ ์ œ์–ด๊ถŒ์ด ๋ช…ํ™•ํ•œ โ€˜์˜ค์ผ€์ŠคํŠธ๋ ˆ์ด์…˜ ํŒจํ„ด(Apache Airflow ๋“ฑ)โ€™์„ ์••๋„์ ์œผ๋กœ ์„ ํ˜ธ

1.2 ์˜ค์ผ€์ŠคํŠธ๋ ˆ์ด์…˜ ์•„ํ‚คํ…์ฒ˜์˜ 4๋Œ€ ํ•ต์‹ฌ ์ปดํฌ๋„ŒํŠธ

  • ํ˜„๋Œ€์ ์ธ ์˜ค์ผ€์ŠคํŠธ๋ ˆ์ด์…˜ ์—”์ง„์€ ๋‚ด๋ถ€์ ์œผ๋กœ ๊ณ ๋„์˜ ๋ถ„์‚ฐ ์‹œ์Šคํ…œ ๊ตฌ์กฐ๋ฅผ ์ฑ„ํƒํ•˜๊ณ  ์žˆ์Œ

    1. ์ปจํŠธ๋กค ํ”Œ๋ ˆ์ธ & ์Šค์ผ€์ค„๋Ÿฌ (Control Plane & Scheduler):
      • ์ „์ฒด ํŒŒ์ดํ”„๋ผ์ธ์˜ ๋ผˆ๋Œ€(DAG)๋ฅผ ํ•ด์„ํ•˜๊ณ ,
      • ๊ฐ ํƒœ์Šคํฌ์˜ ์˜์กด์„ฑ๊ณผ ์ง„์ž…์ฐจ์ˆ˜(In-degree)๋ฅผ ๊ณ„์‚ฐํ•˜์—ฌ
      • ์‹คํ–‰ ๊ฐ€๋Šฅํ•œ ์ž‘์—…์„ ์„ ๋ณ„ํ•˜๋Š” ์—ญํ• 
    2. ๋ฉ”ํƒ€๋ฐ์ดํ„ฐ ์ €์žฅ์†Œ (Metadata Repository):
      • ๋ชจ๋“  ์›Œํฌํ”Œ๋กœ์šฐ์˜ ์‹คํ–‰ ์ด๋ ฅ, ํƒœ์Šคํฌ์˜ ๋ผ์ดํ”„์‚ฌ์ดํด ์ƒํƒœ(Queued, Running, Failed, Success), ์ „์—ญ ์„ค์ •๊ฐ’ ๋“ฑ์„ ์˜๊ตฌ ์ €์žฅํ•˜๋Š” ์•„ํ‚คํ…์ฒ˜์˜ ์‹ฌ์žฅ๋ถ€
      • RDBMS ์ฃผ๋กœ ์‚ฌ์šฉ
    3. ์‹คํ–‰๊ธฐ ๋ฐ ํ ์ธํ”„๋ผ (Executor & Message Queue):
      • ์Šค์ผ€์ค„๋Ÿฌ๊ฐ€ ์„ ๋ณ„ํ•œ ํƒœ์Šคํฌ๋ฅผ ์‹ค์ œ ์ผ๊พผ(Worker)๋“ค์—๊ฒŒ ์•ˆ์ „ํ•˜๊ฒŒ ์ „๋‹ฌํ•˜๊ธฐ ์œ„ํ•œ ๋ฒ„ํผ ๋ฐ ์ค‘๊ณ„ ๊ณ„์ธต
      • Redis, RabbitMQ ๊ฐ™์€ ๋ธŒ๋กœ์ปค๋‚˜ ์ฟ ๋ฒ„๋„คํ‹ฐ์Šค API API์™€ ์—ฐ๋™๋จ
    4. ๋ถ„์‚ฐ ์›Œ์ปค ํด๋Ÿฌ์Šคํ„ฐ (Distributed Workers):
      • ์˜ค์ผ€์ŠคํŠธ๋ ˆ์ดํ„ฐ์˜ ๋ช…๋ น์„ ๋ฐ›์•„ ์‹ค์ œ ์ปดํ“จํŒ… ์—ฐ์‚ฐ(API ํ˜ธ์ถœ, SQL ํŠธ๋ฆฌ๊ฑฐ, Spark ์ง๋ ฌํ™”)์„ ์ˆ˜ํ–‰ํ•˜๋Š” ๋ฌผ๋ฆฌ/๋…ผ๋ฆฌ์  ๋…ธ๋“œ

1.3 ์˜ค์ผ€์ŠคํŠธ๋ ˆ์ด์…˜ ์„ค๊ณ„ ์‹œ ํ•„์ˆ˜ ์•„ํ‚คํ…์ฒ˜ ์›์น™

  • ์„ฑ๊ณต์ ์ธ ์˜ค์ผ€์ŠคํŠธ๋ ˆ์ด์…˜ ํŒŒ์ดํ”„๋ผ์ธ์„ ๊ตฌ์ถ•ํ•˜๊ธฐ ์œ„ํ•ด ์•„ํ‚คํ…์ฒ˜ ๋ ˆ๋ฒจ์—์„œ ๋ฐ˜๋“œ์‹œ ์ค€์ˆ˜ํ•ด์•ผ ํ•˜๋Š” ์—”์ง€๋‹ˆ์–ด๋ง ์›์น™

    • ์ปดํ“จํŒ…๊ณผ ์˜ค์ผ€์ŠคํŠธ๋ ˆ์ด์…˜์˜ ๋ถ„๋ฆฌ (Decoupling)
      • ๊ฐ€์žฅ ์ค‘์š”ํ•œ ์›์น™ ๐Ÿกฒ โ€œ์˜ค์ผ€์ŠคํŠธ๋ ˆ์ดํ„ฐ๋Š” ์‹ ํ˜ธ๋“ฑ ์—ญํ• ๋งŒ ํ•ด์•ผ์ง€, ์ง์ ‘ ์ฐจ๊ฐ€ ๋˜์–ด์„œ๋Š” ์•ˆ ๋œ๋‹คโ€๋Š” ๋ฒ•์น™
      • ์˜ค์ผ€์ŠคํŠธ๋ ˆ์ดํ„ฐ ์›Œ์ปค ์ž์ฒด์˜ ๋ฉ”๋ชจ๋ฆฌ์™€ CPU๋ฅผ ์†Œ๋ชจํ•˜์—ฌ ๋Œ€์šฉ๋Ÿ‰ ๋ฐ์ดํ„ฐ๋ฅผ ์ฒ˜๋ฆฌ(์˜ˆ: Pandas ๋ฐ์ดํ„ฐ ๋ณ€ํ™˜)ํ•˜๋ฉด ์‹œ์Šคํ…œ ์ „์ฒด๊ฐ€ ๋งˆ๋น„๋จ
      • ๋ฌด๊ฑฐ์šด ์—ฐ์‚ฐ์€ ๋ถ„์‚ฐ ์ปดํ“จํŒ… ์—”์ง„(Spark, Ray, Trino)์ด๋‚˜ ์™ธ๋ถ€ DB ์—”์ง„์— ์œ„์ž„ํ•˜๊ณ ,
      • ์˜ค์ผ€์ŠคํŠธ๋ ˆ์ดํ„ฐ๋Š” ์‹คํ–‰ ๋ช…๋ น(Trigger)๊ณผ ์™„๋ฃŒ ์—ฌ๋ถ€ ํ™•์ธ(Polling)๋งŒ ์ˆ˜ํ–‰ํ•ด์•ผ ํ•จ
    • ๋ฉฑ๋“ฑ์„ฑ (Idempotency) ์ธํ”„๋ผ ๊ตฌ์ถ•
      • ์˜ค์ผ€์ŠคํŠธ๋ ˆ์ด์…˜ ์•„ํ‚คํ…์ฒ˜์—์„œ๋Š” ํŠน์ • ํƒœ์Šคํฌ๊ฐ€ ์‹คํŒจํ•˜์—ฌ ์žฌ์‹คํ–‰(Retry)๋˜๊ฑฐ๋‚˜ ๊ณผ๊ฑฐ ํŠน์ • ์‹œ์ ์œผ๋กœ ๋Œ์•„๊ฐ€ ๋ฐฑํ•„(Backfill)์„ ์ˆ˜ํ–‰ํ•  ๋•Œ,
        • ๋ช‡ ๋ฒˆ์„ ๋‹ค์‹œ ์‹คํ–‰ํ•ด๋„ ํƒ€๊ฒŸ ์ €์žฅ์†Œ์˜ ์ตœ์ข… ๋ฐ์ดํ„ฐ์…‹ ๊ฒฐ๊ณผ๊ฐ€ ํ•ญ์ƒ ๋™์ผํ•จ์„ ๋ณด์žฅํ•ด์•ผ ํ•จ
      • ๋ฐ์ดํ„ฐ ์†Œ์Šค๋ฅผ ๊ฒฉ๋ฆฌํ•  ์ˆ˜ ์žˆ๋Š” ๋…ผ๋ฆฌ์  ์‹œ์  ๋ณ€์ˆ˜(Logical Date) ๋ฐ”์ธ๋”ฉ ๋ฐ ์ €์žฅ์†Œ์˜ Upsert ๋ฉ”์ปค๋‹ˆ์ฆ˜ ์„ค๊ณ„๊ฐ€ ์•„ํ‚คํ…์ฒ˜์— ๋‚ด์žฌ๋˜์–ด์•ผ ํ•จ
    • ๊ฒฐํ•ฉ ๊ฒฉ๋ฆฌ ๋ฐ ์ž์› ๊ฒฉ๋ฆฌ (Isolation)
      • ํŒŒ์ดํ”„๋ผ์ธ ๋‚ด์˜ ๊ฐ ๋‹จ๊ณ„๋Š” ์ƒํ˜ธ ๊ฐ„์˜ ๋ผ์ด๋ธŒ๋Ÿฌ๋ฆฌ ์˜์กด์„ฑ์ด๋‚˜ ํ•˜๋“œ์›จ์–ด ์ž์› ์†Œ๋น„์— ์˜ํ–ฅ์„ ์ฃผ์ง€ ์•Š์•„์•ผ ํ•จ
      • AI ํŒŒ์ดํ”„๋ผ์ธ์—์„œ๋Š” ์ผ๋ฐ˜ ๊ฐ€๋ฒผ์šด SQL ์ „์ฒ˜๋ฆฌ ํƒœ์Šคํฌ์™€ ๋Œ€์šฉ๋Ÿ‰ GPU ๊ฐ€์†์ด ํ•„์š”ํ•œ ๋ชจ๋ธ ํ•™์Šต ํƒœ์Šคํฌ๊ฐ€ ๊ณต์กดํ•˜๋ฏ€๋กœ,
      • ํƒœ์Šคํฌ ๋‹จ์œ„๋ฅผ ์ปจํ…Œ์ด๋„ˆ(Docker Pod) ํ˜•ํƒœ๋กœ ๋™์  ๊ฒฉ๋ฆฌํ•˜์—ฌ
      • ํ•„์š”ํ•œ ์ธํ”„๋ผ์— ์‹ค์‹œ๊ฐ„ ๋ฐฐ์ •ํ•˜๋Š” ์•„ํ‚คํ…์ฒ˜(์˜ˆ: KubernetesPodOperator)๋ฅผ ๊ตฌ์ถ•ํ•˜๋Š” ๊ฒƒ์ด ์ตœ์„ 

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๋‹จ๊ณ„ ํ”„๋กœ์„ธ์Šค
  1. [Task 1] Ingestion & Lake Landing (์ˆ˜์ง‘ ๋ฐ ์ ์žฌ):
    1. ์™ธ๋ถ€ ๋ฐ์ดํ„ฐ๋ฅผ ๋‹ค์šด๋กœ๋“œํ•˜์—ฌ
    2. MinIO์˜ raw-data/ ๋ฒ„ํ‚ท์— ์ €์žฅ
  2. [Task 2] Spark ๋ถ„์‚ฐ ์—ฐ์‚ฐ ํŠธ๋ฆฌ๊ฑฐ (Spark-Submit):
    1. Airflow๊ฐ€ Spark Operator๋ฅผ ํ†ตํ•ด ์ „์ฒ˜๋ฆฌ ์ž‘์—…์„ ๋ช…๋ น
    2. Spark ํด๋Ÿฌ์Šคํ„ฐ๊ฐ€ MinIO์˜ ์›์‹œ ๋ฐ์ดํ„ฐ๋ฅผ ์ฝ์–ด์™€ ์ฒญํ‚น(Chunking)์„ ์ˆ˜ํ–‰
    3. processed-data/ ๋ฒ„ํ‚ท์— Parquet ํ˜•์‹์œผ๋กœ ์ €์žฅ
  3. [Task 3] ๋ณ‘๋ ฌ ๋ถ„์‚ฐ ์ž„๋ฒ ๋”ฉ ๋ฐ Vector DB Upsert:
    1. ๊ฐ€๊ณต๋œ ํ…์ŠคํŠธ ์ฒญํฌ๋“ค์„ ์ฝ์–ด์™€
    2. ๋กœ์ปฌ AI ์ž„๋ฒ ๋”ฉ ๋ชจ๋ธ(Ollama/HuggingFace)์„ ํ†ตํ•ด ๊ณ ์ฐจ์› ๋ฒกํ„ฐ๋กœ ๋ณ€ํ™˜ํ•œ ๋’ค,
    3. Qdrant์˜ HNSW ๊ทธ๋ž˜ํ”„ ์ธ๋ฑ์Šค์— Upsert
  4. [Task 4] ์ธ๋ฑ์Šค ์ •๋น„ ๋ฐ ์บ์‹œ ํด๋ฆฐ์—…:
    • ๋ฉ”ํƒ€๋ฐ์ดํ„ฐ ๊ฐฑ์‹  ๋ฐ ๋ฆฌ์†Œ์Šค ํ•ด์ œ

2. ์ข…ํ•ฉ ์‹ค์Šต ์˜ˆ์ œ ์ฝ”๋“œ 1

  • ํ˜„๋Œ€์ ์ธ AI ์ธํ”„๋ผ์˜ ํ‘œ์ค€ ๊ตฌ์กฐ์ธ ํ•˜์ด๋ธŒ๋ฆฌ๋“œ RAG(๊ฒ€์ƒ‰ ์ฆ๊ฐ• ์ƒ์„ฑ) ํ”Œ๋žซํผ์˜ ๋ฐ์ดํ„ฐ ๊ณต๊ธ‰์„ ์„ ์•„ํ‚คํ…์ฒ˜ ๊ด€์ ์—์„œ ํ”„๋กœํ† ํƒ€์ดํ•‘ํ•œ ํ•ต์‹ฌ ์„ค๊ณ„๋„
  • โ€œ์ˆ˜์ง‘ ๐Ÿกฒ Data Lake ๐Ÿกฒ Apache Spark ๐Ÿกฒ VectorDBโ€ ๊ตฌ์กฐ๋ฅผ ๋‹จ์ผ ํŒŒ์ผ๋กœ ๊ตฌํ˜„ํ•œ Airflow DAG
    • ์‹ค๋ฌด์—์„œ๋Š” Spark ์ „์ฒ˜๋ฆฌ ๋กœ์ง์„ ๋ณ„๋„์˜ *.py ํŒŒ์ผ๋กœ ๋ถ„๋ฆฌํ•˜์—ฌ SparkSubmitOperator๋กœ ํ˜ธ์ถœ
    • ์‹ค์Šต์˜ˆ์ œ์—์„œ๋Š” ๊ฐ€๋…์„ฑ์„ ์œ„ํ•ด PySpark ์ „์ฒ˜๋ฆฌ ๋ฐ Qdrant ์ ์žฌ๋ฅผ ํ†ตํ•ฉ ๊ตฌํ˜„ํ•จ

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_exists API๋ฅผ ํ†ตํ•ด ๋ฐฉ์–ด์  ์ฝ”๋“œ(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 ์ธ๋ฑ์‹ฑ)
    
    1. ์ง„์ž… ๊ด€๋ฌธ ๋ฐ ๋ฐ์ดํ„ฐ ๋ ˆ์ดํฌ ์•ˆ์ฐฉ (task_ingest):
      • Airflow ์Šค์ผ€์ค„๋Ÿฌ๊ฐ€ ์ง„์ž…์ฐจ์ˆ˜(In-degree=0)๊ฐ€ ์ œ๋กœ์ธ task_ingest๋ฅผ ๊ฐ€์žฅ ๋จผ์ € ํ์— ๋„ฃ๊ณ  ์›Œ์ปค์— ๋ฐฐ์ •
      • ๊ฐ€์ƒ์˜ ๋น„์ •ํ˜• ์†Œ์Šค ํ…์ŠคํŠธ ๋ฐ์ดํ„ฐ๊ฐ€
        • ๋ฐ”์ดํŠธ ์ŠคํŠธ๋ฆผ(io.BytesIO) ํ˜•ํƒœ๋กœ ๋ณ€ํ™˜๋˜์–ด
        • ๋„คํŠธ์›Œํฌ ๋ง์„ ํƒ€๊ณ 
        • MinIO ์˜ค๋ธŒ์ ํŠธ ์Šคํ† ๋ฆฌ์ง€ ๋‚ด๋ถ€์˜ raw/manual_01.txt ๊ฒฝ๋กœ๋กœ ์—…๋กœ๋“œ ๋ฐ ๊ฒฉ๋ฆฌ
    2. ์ˆœ์ฐจ ์ œ์–ด๊ถŒ ์ดํ–‰ (>>):
      • 1๋‹จ๊ณ„ ํƒœ์Šคํฌ๊ฐ€ ์„ฑ๊ณต(Success)์œผ๋กœ ๋งˆํ‚น๋˜๋ฉด,
      • ์˜์กด์„ฑ ๊ฒฐํ•ฉ ์—ฐ์‚ฐ์ž(>>)๋ฅผ ํƒ€๊ณ 
      • ์ œ์–ด๊ถŒ์ด ๋‹ค์Œ ํƒœ์Šคํฌ์ธ task_spark_and_vector๋กœ ์•ˆ์ „ํ•˜๊ฒŒ ์ „์ด
    3. ๋ฐ์ดํ„ฐํ•˜์šฐ์Šค ๋กœ๋“œ ๋ฐ ๋ฉ”๋ชจ๋ฆฌ ์ฒญํ‚น ์—ฐ์‚ฐ:
      • ๋‘ ๋ฒˆ์งธ ํƒœ์Šคํฌ๊ฐ€ ๊ธฐ๋™๋˜๋ฉด์„œ ๋ฐฉ๊ธˆ MinIO์— ๋ฐฑ์—…๋˜์—ˆ๋˜ ์›์‹œ ํŒŒ์ผ์˜ ์Šค๋ƒ…์ƒท์„ ๋ฉ”๋ชจ๋ฆฌ๋กœ ๋‹ค์‹œ ๋‹ค์šด๋กœ๋“œ
      • replace ๊ฐ€๊ณต์„ ํ†ตํ•ด ๋ฌธ์ž์—ด ๋…ธ์ด์ฆˆ๋ฅผ ์ •์ œํ•˜๊ณ ,
        • ์Šฌ๋ผ์ด๋“œ ์œˆ๋„์šฐ ๋ฐฉ์‹์œผ๋กœ 50๊ธ€์ž์”ฉ ์ชผ๊ฐœ์ง„ ํ…์ŠคํŠธ ์ŠคํŠธ๋ง ๋ฆฌ์ŠคํŠธ(chunks)๋ฅผ ์ƒ์„ฑํ•˜์—ฌ
        • ๋ฉ”๋ชจ๋ฆฌ์— ๋ถ„์‚ฐ ๋ฐฐ์น˜
    4. ๋ฉฑ๋“ฑ์„ฑ ๊ธฐ๋ฐ˜ ๊ณ ์ฐจ์› ๋ฒกํ„ฐ์Šคํ† ์–ด ์ตœ์ข… ๋™๊ธฐํ™”:
      • Qdrant ํด๋ผ์ด์–ธํŠธ๊ฐ€ ๊ฐ€๋™๋˜์–ด
        • ๋ฒกํ„ฐ์Šคํ† ์–ด ๋‚ด๋ถ€์— Distance.COSINE ๊ฑฐ๋ฆฌ๋ฅผ ์—ฐ์‚ฐํ•  ์ˆ˜ ์žˆ๋Š” ์ธ๋ฑ์Šค ๋ ˆ์ด์–ด๋ฅผ ์„ ์–ธ
      • ์ฒญํ‚น๋œ ๋ฌธ์ž์—ด ๋ฐ์ดํ„ฐ ๊ฐ๊ฐ์— ์œ ์ผํ•œ ID ๊ณ ์œ  ๋ฒˆํ˜ธ(idx + 100)๋ฅผ ๊ฐ•์ œ๋กœ ๋ถ€์—ฌ
      • ์ด๋ ‡๊ฒŒ ๊ณ ์œ  ID๋ฅผ ๋ฐ”์ธ๋”ฉํ•จ์œผ๋กœ์จ,
        • ์ด ํŒŒ์ดํ”„๋ผ์ธ์„ ํ•˜๋ฃจ์— ์ˆ˜์‹ญ ๋ฒˆ ์ค‘๋ณตํ•ด์„œ ๋‹ค์‹œ ์‹คํ–‰ํ•˜๋”๋ผ๋„
        • Qdrant ๋‚ด๋ถ€์— ๋ฐ์ดํ„ฐ๊ฐ€ ๋ˆ„์ ๋˜์ง€ ์•Š๊ณ  ๋ฎ์–ด์จ์ง€๊ฒŒ(Upsert) ๋งŒ๋“ค์–ด
        • ๋ฐ์ดํ„ฐ์˜ ๋ฉฑ๋“ฑ์„ฑ(Idempotency)์„ ์ตœ์ข… ์™„์ˆ˜ํ•˜๊ณ 
        • ํŒŒ์ดํ”„๋ผ์ธ ์ „์ฒด๊ฐ€ ์ •์ƒ ์ข…๋ฃŒ๋จ

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
        • ์ปจํ…Œ์ด๋„ˆ์˜ ๋ฒ„์ „์„ ์ˆ˜์ •ํ•˜๊ฑฐ๋‚˜ ํŒŒ์ด๋ธŒ๋Ÿฌ๋ฆฌ ์„ค์น˜ ์š”์ฒญ์—์„œ ๋ฒ„์ „์„ ์ง€์ •ํ•  ๊ฒƒ
    • ์˜ˆ์ œ 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())
          


  • 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 ์ง€์‹ ๋ฒ ์ด์Šค ๋™๊ธฐํ™”โ€

  1. ์‹ค์‹œ๊ฐ„ ์ˆ˜์ง‘ (Kafka):
    • ๊ณต์žฅ ์„ผ์„œ ๋ฐ ์„ค๋น„ ์ œ์–ด๊ธฐ์—์„œ ๋ฐœ์ƒํ•˜๋Š” ์‹ค์‹œ๊ฐ„ ๋น„์ •ํ˜• ๋กœ๊ทธ ๋ฐ ๋งค๋‰ด์–ผ ํ…์ŠคํŠธ๊ฐ€
    • Kafka ํ† ํ”ฝ์œผ๋กœ ์ธ๋ฑ์‹ฑ
  2. ๋ฐ์ดํ„ฐ ๋ ˆ์ดํฌ ๋ณด์กด (MinIO):
    • ์ˆ˜์ง‘๋œ ์›์‹œ(Raw) ๋กœ๊ทธ๋Š” ๋ฐ์ดํ„ฐ ์œ ์‹ค ๋ฐฉ์ง€ ๋ฐ ๋ฉฑ๋“ฑ์„ฑ ํ™•๋ณด๋ฅผ ์œ„ํ•ด
    • MinIO ์˜ค๋ธŒ์ ํŠธ ์Šคํ† ๋ฆฌ์ง€์˜ ๋‚ ์งœ๋ณ„ ํŒŒํ‹ฐ์…˜ ์˜์—ญ(raw/{ { ds_nodash }}/)์— ์˜๊ตฌ ๋ณด์กด
  3. ๋ถ„์‚ฐ ์ „์ฒ˜๋ฆฌ ๋ฐ ์ž„๋ฒ ๋”ฉ ๊ฐ€๊ณต (PySpark):
    • Spark ์„ธ์…˜์„ ์ปจํ…Œ์ด๋„ˆ ๋‚ด๋ถ€์—์„œ ๋™์  ๊ตฌ๋™ํ•˜์—ฌ,
    • ๋น„์ •ํ˜• ํ…์ŠคํŠธ์˜ ๋…ธ์ด์ฆˆ๋ฅผ ์ œ๊ฑฐํ•˜๊ณ 
    • ์˜๋ฏธ๋ก ์  ๋ฌธ๋งฅ ๋ณด์ „์„ ์œ„ํ•œ ์ฒญํ‚น(Chunking) ๋ถ„์‚ฐ ์—ฐ์‚ฐ์„ ์ˆ˜ํ–‰
  4. ์ง€์‹ ๋ฒ ์ด์Šค ๋™๊ธฐํ™” (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
      


  • 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:
      

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

3.5 DAG ๊ตฌ์„ฑ

  • ํŒŒ์ดํ”„๋ผ์ธ ์˜์กด์„ฑ ์•„ํ‚คํ…์ฒ˜ ๋ฐ ํ๋ฆ„๋„
    • ๊ฐ ์˜คํ”ˆ์†Œ์Šค ์ปดํฌ๋„ŒํŠธ์˜ ๊ฒฐํ•จ ๊ฒฉ๋ฆฌ(Fault Isolation)๋ฅผ ์‹œ๊ฐํ™”ํ•œ ์œ„์ƒ ์ •๋ ฌ ์˜์กด์„ฑ ๊ตฌ์กฐ


  • ๋ฐ์ดํ„ฐ ์ œ์–ด ํ๋ฆ„ ๊ตฌ์กฐ (Data Control Flow)

    • ๋‹จ๊ณ„๋ณ„ ๋งค์ปค๋‹ˆ์ฆ˜
      1. task_kafka_kr_ingest (์ˆ˜์ง‘ ๊ณ„์ธต):
        • KRaft ๋‹จ๋… ๋…ธ๋“œ๋กœ ๊ธฐ๋™ ์ค‘์ธ Kafka ํ† ํ”ฝ์œผ๋กœ
        • ๊ฐ€๋ณ€ ๋ฉด์  ์œ ๋Ÿ‰๊ณ„ ์žฅ์•  ๋กœ๊ทธ ์ด๋ฒคํŠธ๋ฅผ ๋ฐœํ–‰ํ•˜๊ณ  ์ˆ˜์ง‘ ๊ฒ€์ฆ
      2. ๋ณ‘๋ ฌ ์ฒ˜๋ฆฌ ๊ณ„์ธต (Parallel Operations):
        • ์ˆ˜์ง‘ ์™„๋ฃŒ ์‹ ํ˜ธ๋ฅผ ๋ฐ›์œผ๋ฉด,
        • ๋ฐ์ดํ„ฐ ๋ ˆ์ดํฌ ๋ฐฑ์—…(MinIO), ๋ถ„์‚ฐ ๋ฉ”๋ชจ๋ฆฌ ์ „์ฒ˜๋ฆฌ(Spark)๊ฐ€
        • ํ˜ธ์ŠคํŠธ ์ž์›์„ ํšจ์œจ์ ์œผ๋กœ ๋‚˜๋ˆ„์–ด ์“ฐ๋ฉฐ ๋™์‹œ์— ๋ณ‘๋ ฌ๋กœ ๊ธฐ๋™
      3. 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_v5 DAG๊ฐ€ ๋‚˜ํƒ€๋‚˜๋Š”์ง€ ๊ฒ€์ƒ‰์ฐฝ์—์„œ ์กฐํšŒ

  • ๊ฒฐ๊ณผ ํ™•์ธ

    1. ํ™•์ธ ์ „ ์ƒํ™ฉ
      • ๊ธฐ์กด์— ์‹คํŒจํ–ˆ๋˜ ์„ค์ • ์ดํ›„๋กœ ์„ฑ๊ณตํ•œ ์„ค์ •์„ ํ™•์ธํ•  ์ˆ˜ ์žˆ์Œ
      • ์ดˆ๋ฐ˜์˜ ์„ฑ๊ณต ์ดํ›„, ์„ค์ •์ด ์•ˆ์ •ํ™”๋œ ๋’ค์—๋Š” ์ฒ˜๋ฆฌ ์‹œ๊ฐ„๋„ ํฌ๊ฒŒ ๊ฐ์†Œํ•˜์˜€์Œ
    2. Trigger ์‹คํ–‰ ๊ฒฐ๊ณผ
      • ์ „ ๊ณผ์ •์ด ๋ฌด๋ฆฌ์—†์ด ์„ฑ๊ณตํ•จ
      • ์ฒ˜๋ฆฌ ์‹œ๊ฐ„: ์‹œ์ž‘๋ถ€ํ„ฐ ๋๊นŒ์ง€ ์•ฝ 8์ดˆ์— ์™„๋ฃŒ๋˜์—ˆ์Œ
        • ์„ค์ • ์ˆ˜์ •์— ๋”ฐ๋ฅธ ์•ˆ์ •ํ™” ํ›„ ์†Œ์š” ์‹œ๊ฐ„์ด ๊ธ‰๊ฐํ–ˆ์Œ์„ ํ™•์ธ
    3. task_kafka_kr_ingest์˜ ๋กœ๊ทธ
      • 3๊ฐœ์˜ Kafka ๋ณผ๋ฅจ์— ์ œ๋Œ€๋กœ ๋ณต์ œ, ์ €์žฅ๋˜์—ˆ์Œ
    4. task_spark_transform์˜ ๋กœ๊ทธ
      • ๊ธฐ์กด์— ์‹คํŒจํ–ˆ๋˜ ์„ค์ • ์ดํ›„๋กœ ์„ฑ๊ณตํ•œ ์„ค์ •์„ ํ™•์ธํ•  ์ˆ˜ ์žˆ์Œ
      • ์ค‘๊ฐ„์— Error ํ‘œ์‹œ๊ฐ€ ๋ณด์ด์ง€๋งŒ ๋‚ด์šฉ์„ ์ฝ์–ด๋ณด๋ฉด Error๊ฐ€ ์•„๋‹˜์„ ์•Œ ์ˆ˜ ์žˆ์Œ
        • ERROR ํ‘œ์‹œ๊ฐ€ ์ถœ๋ ฅ๋œ ๊ฒƒ์€ ์ŠคํŒŒํฌ์˜ ํ‘œ์ค€ ์ž๋ฐ” ์—๋Ÿฌ ์ถœ๋ ฅ(Stderr) ์ŠคํŠธ๋ฆผ์ด ์œ ์ž…๋˜์–ด ๋ฐœ์ƒํ•œ ์—์–ดํ”Œ๋กœ์šฐ ๊ณ ์œ ์˜ ๋กœ๊น… ํŠน์ง•
          • ์ŠคํŒŒํฌ ์—”์ง„์€ ๊ตฌ๋™๋  ๋•Œ ์—”์ง„์˜ ์‹œ์Šคํ…œ ๊ฒฝ๊ณ ๋‚˜ ํ™˜๊ฒฝ ์„ค์ • ์•ˆ๋‚ด๋ฅผ Stdout์ด ์•„๋‹Œ Stderr ์ŠคํŠธ๋ฆผ์œผ๋กœ ๋‚ด๋ณด๋‚ด๋„๋ก ์ฝ”๋“œ๊ฐ€ ์ž‘์„ฑ๋˜์–ด ์žˆ์Œ
          • ์—์–ดํ”Œ๋กœ์šฐ๊ฐ€ ์ถœ๋ ฅ ์ŠคํŠธ๋ฆผ์ด Stderr๋กœ ์ „๋‹ฌ๋œ ๊ฒƒ๋งŒ ๋ณด๊ณ  ๊ธฐ๊ณ„์ ์œผ๋กœ ERROR๋ฅผ ๋ถ™์—ฌ๋ฒ„๋ฆฐ ๊ฒƒ์ž„
        • ์ง„์งœ ๋‚ด๋ถ€ ์˜ˆ์™ธ(Exception)๊ฐ€ ๋ฐœ์ƒํ–ˆ๋‹ค๋ฉด ์ŠคํŒŒํฌ ํŠน์œ ์˜ ๊ฑฐ๋Œ€ํ•œ StackTrace ์ž๋ฐ” ์—๋Ÿฌ ๋ฌธ๋‹จ์ด ์ถœ๋ ฅ๋˜์–ด์•ผ ํ•จ
    5. task_minio_raw_backup ๋กœ๊ทธ
      • ๋ฐ์ดํ„ฐ์˜ ์ ์žฌ๊ฐ€ ์ •์ƒ์ ์œผ๋กœ ์™„๋ฃŒ๋˜์—ˆ์Œ
    6. task_vector_upsert_qdrant์˜ ๋กœ๊ทธ
      • ๋ฐ์ดํ„ฐ์˜ Upsert๊ฐ€ ์ •์ƒ์ ์œผ๋กœ ์™„๋ฃŒ๋˜์—ˆ์Œ

๐Ÿ’ฌ Q&A ๋ฐ ์˜๊ฒฌ ๋‚˜๋ˆ„๊ธฐ