<?xml version="1.0" encoding="UTF-8"?>
<rss version="2.0" xmlns:atom="http://www.w3.org/2005/Atom" xmlns:dc="http://purl.org/dc/elements/1.1/">
  <channel>
    <title>DEV Community: 박준현</title>
    <description>The latest articles on DEV Community by 박준현 (@junhyun-dev).</description>
    <link>https://dev.to/junhyun-dev</link>
    <image>
      <url>https://media2.dev.to/dynamic/image/width=90,height=90,fit=cover,gravity=auto,format=auto/https:%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Fuser%2Fprofile_image%2F4021062%2Fcffaa152-163d-4a4c-88d8-8f2e9c688b0e.png</url>
      <title>DEV Community: 박준현</title>
      <link>https://dev.to/junhyun-dev</link>
    </image>
    <atom:link rel="self" type="application/rss+xml" href="https://dev.to/feed/junhyun-dev"/>
    <language>en</language>
    <item>
      <title>Python batch를 Spark로 옮겼더니 802.675가 달라졌다</title>
      <dc:creator>박준현</dc:creator>
      <pubDate>Sun, 19 Jul 2026 15:23:47 +0000</pubDate>
      <link>https://dev.to/junhyun-dev/python-batchreul-sparkro-olmgyeossdeoni-802675ga-dalrajyeossda-3li0</link>
      <guid>https://dev.to/junhyun-dev/python-batchreul-sparkro-olmgyeossdeoni-802675ga-dalrajyeossda-3li0</guid>
      <description>&lt;h1&gt;
  
  
  Python batch를 Spark로 옮겼더니 802.675가 달라졌다
&lt;/h1&gt;

&lt;p&gt;Python으로 만든 batch transform을 Spark DataFrame으로 다시 쓰는 일은 문법 번역처럼 보인다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;Python loop / dict
-&amp;gt; Spark select / Window / groupBy
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;하지만 실행 엔진만 바꿨는데도 결과가 달라질 수 있다. 중복 중 어떤 행을 남기는지, 평균을 어떻게 반올림하는지, quality 실패가 orchestration에 어떻게 전달되는지, 같은 입력을 정말 skip해도 되는지가 모두 계약이기 때문이다.&lt;/p&gt;

&lt;p&gt;&lt;code&gt;manufacturing-data-platform-mini&lt;/code&gt;의 S7에서는 Kafka landing에서 만든 한 날짜의 canonical CSV를 기존 Python batch와 새 Spark batch에 함께 넣었다. 목표는 "Spark를 사용했다"가 아니었다.&lt;/p&gt;

&lt;blockquote&gt;
&lt;p&gt;같은 입력을 다른 엔진으로 처리해도 grain, dedup, metric, quality, publish 상태 전이가 같아야 한다.&lt;/p&gt;
&lt;/blockquote&gt;

&lt;p&gt;그 계약을 테스트하던 중 실제로 &lt;code&gt;802.675&lt;/code&gt; 반올림 결과가 갈리는 false-green도 잡았다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Scenario
&lt;/h2&gt;

&lt;p&gt;앞선 slice에서는 local Kafka에서 받은 합성 설비 event를 immutable JSONL로 landing한 뒤, 한 &lt;code&gt;business_date&lt;/code&gt;의 accepted event를 canonical CSV와 &lt;code&gt;source_hash&lt;/code&gt;로 고정했다.&lt;/p&gt;

&lt;p&gt;S7은 그 경계부터 시작한다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;Kafka raw landing
-&amp;gt; K1.5 canonical CSV + source_hash
-&amp;gt; Spark silver / gold
-&amp;gt; existing quality checks
-&amp;gt; Iceberg business_date partition
-&amp;gt; thin Airflow CLI wrapper
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;운영 시나리오는 작다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;한 날짜를 Spark batch로 backfill한다.
같은 source 재실행은 새 snapshot을 만들지 않는다.
정정 source는 대상 날짜 partition만 교체한다.
quality가 실패하면 Iceberg current를 바꾸지 않는다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Spark가 Kafka JSONL을 다시 해석하게 만들지 않았다. 이미 provenance와 형식을 검증한 adapter의 CSV와 &lt;code&gt;source_hash&lt;/code&gt;를 그대로 입력 계약으로 재사용했다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Decision Pressure: 엔진보다 먼저 고정할 것
&lt;/h2&gt;

&lt;p&gt;Python 구현과 Spark 구현이 "같다"고 하려면 무엇을 비교해야 할까?&lt;/p&gt;

&lt;p&gt;이번 slice에서는 다섯 가지를 먼저 고정했다.&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;계약&lt;/th&gt;
&lt;th&gt;고정한 내용&lt;/th&gt;
&lt;th&gt;깨지면 생기는 문제&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;input identity&lt;/td&gt;
&lt;td&gt;canonical CSV + &lt;code&gt;source_hash&lt;/code&gt;
&lt;/td&gt;
&lt;td&gt;서로 다른 입력을 같은 결과처럼 비교&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;grain / dedup&lt;/td&gt;
&lt;td&gt;natural key와 Kafka coordinate 기준 first&lt;/td&gt;
&lt;td&gt;합계 중복 또는 임의 행 선택&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;metric rounding&lt;/td&gt;
&lt;td&gt;Python gold와 같은 소수 둘째 자리 결과&lt;/td&gt;
&lt;td&gt;엔진별 metric 불일치&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;quality gate&lt;/td&gt;
&lt;td&gt;기존 quality suite를 그대로 적용&lt;/td&gt;
&lt;td&gt;Spark 경로만 다른 품질 기준 사용&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;publish state&lt;/td&gt;
&lt;td&gt;same-source skip, correction overwrite&lt;/td&gt;
&lt;td&gt;snapshot noise 또는 stale skip&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;이것이 코드보다 먼저 필요했다. &lt;code&gt;groupBy&lt;/code&gt;가 실행된다는 사실만으로는 기존 batch와 같은 시스템이 되지 않는다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Contract 1: 검증된 입력 경계를 재사용한다
&lt;/h2&gt;

&lt;p&gt;S7은 K1.5 adapter가 만든 결과만 받는다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;input rows   = canonical CSV
input version = SHA-256 source_hash
date scope    = exactly one business_date
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Spark가 raw Kafka envelope을 직접 읽게 하면 adapter에서 검증한 manifest, accepted status, coordinate, provenance 계약을 우회한다. 엔진 교체 범위를 줄이기 위해 source boundary는 바꾸지 않았다.&lt;/p&gt;

&lt;p&gt;또한 요청한 날짜와 CSV 내부 날짜가 다르면 SparkSession을 띄우거나 table을 건드리기 전에 실패시킨다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Contract 2: dedup의 "첫 행"도 정의한다
&lt;/h2&gt;

&lt;p&gt;silver natural key는 다음과 같다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;(work_order_id, machine_id, event_time)
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Python 구현은 canonical CSV에서 먼저 나온 행을 남긴다. 그런데 Spark DataFrame에는 별도의 순서를 주지 않으면 "첫 행"이 결정적이지 않다.&lt;/p&gt;

&lt;p&gt;adapter가 CSV를 Kafka coordinate 순으로 쓰므로 Spark도 같은 순서를 명시했다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight python"&gt;&lt;code&gt;&lt;span class="n"&gt;dedup_order&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;Window&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;partitionBy&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;
    &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;work_order_id&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;machine_id&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;event_time&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;
&lt;span class="p"&gt;).&lt;/span&gt;&lt;span class="nf"&gt;orderBy&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;
    &lt;span class="n"&gt;F&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;col&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;kafka_topic&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;),&lt;/span&gt;
    &lt;span class="n"&gt;F&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;col&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;kafka_partition&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;).&lt;/span&gt;&lt;span class="nf"&gt;cast&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;long&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;),&lt;/span&gt;
    &lt;span class="n"&gt;F&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;col&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;kafka_offset&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;).&lt;/span&gt;&lt;span class="nf"&gt;cast&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;long&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;),&lt;/span&gt;
&lt;span class="p"&gt;)&lt;/span&gt;

&lt;span class="n"&gt;deduped&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="p"&gt;(&lt;/span&gt;
    &lt;span class="n"&gt;filtered&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;withColumn&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;_rn&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;F&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;row_number&lt;/span&gt;&lt;span class="p"&gt;().&lt;/span&gt;&lt;span class="nf"&gt;over&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;dedup_order&lt;/span&gt;&lt;span class="p"&gt;))&lt;/span&gt;
    &lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;filter&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;F&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;col&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;_rn&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt; &lt;span class="o"&gt;==&lt;/span&gt; &lt;span class="mi"&gt;1&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
    &lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;drop&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;_rn&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
&lt;span class="p"&gt;)&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;같은 natural key를 가진 두 event를 넣는 테스트에서 Python과 Spark가 같은 행 하나를 남기는지 확인했다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Contract 3: &lt;code&gt;round&lt;/code&gt; 이름이 같아도 결과는 같지 않다
&lt;/h2&gt;

&lt;p&gt;처음에는 Python &lt;code&gt;round&lt;/code&gt;가 bankers' rounding을 쓰므로 Spark &lt;code&gt;bround&lt;/code&gt;와 같을 것이라고 생각했다. 일반 샘플은 통과했다.&lt;/p&gt;

&lt;p&gt;반례는 &lt;code&gt;cycle_time_ms&lt;/code&gt; 합계 &lt;code&gt;32107&lt;/code&gt;, 행 수 &lt;code&gt;40&lt;/code&gt;인 평균이었다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;수학적 평균: 32107 / 40 = 802.675

Python round(value, 2): 802.67
Spark bround(value, 2): 802.68
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이 차이는 "둘 다 half-even"이라는 이름만 비교해서는 잡히지 않았다. 실행 중인 double 값과 각 함수의 처리 경로까지 포함해야 했다.&lt;/p&gt;

&lt;p&gt;40,400개의 이 프로젝트용 정수비 표본을 비교했을 때 &lt;code&gt;bround&lt;/code&gt;는 Python 결과와 204건 달랐다. &lt;code&gt;format_number&lt;/code&gt;로 둘째 자리까지 만든 뒤 grouping comma를 제거하고 double로 바꾸는 built-in 표현은 같은 표본에서 mismatch가 없었다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight python"&gt;&lt;code&gt;&lt;span class="k"&gt;def&lt;/span&gt; &lt;span class="nf"&gt;_round_like_python&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;col&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;scale&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="nb"&gt;int&lt;/span&gt;&lt;span class="p"&gt;):&lt;/span&gt;
    &lt;span class="k"&gt;return&lt;/span&gt; &lt;span class="n"&gt;F&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;regexp_replace&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;
        &lt;span class="n"&gt;F&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;format_number&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;col&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;scale&lt;/span&gt;&lt;span class="p"&gt;),&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;,&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="sh"&gt;""&lt;/span&gt;
    &lt;span class="p"&gt;).&lt;/span&gt;&lt;span class="nf"&gt;cast&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;double&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;그리고 &lt;code&gt;32107 / 40&lt;/code&gt;을 별도 golden test로 남겼다.&lt;/p&gt;

&lt;p&gt;중요한 경계도 있다.&lt;/p&gt;

&lt;blockquote&gt;
&lt;p&gt;이것은 이 gold metric과 40,400개 정수비 표본에서 확인한 bounded parity다. 모든 double과 모든 scale에서 Python &lt;code&gt;round&lt;/code&gt;와 보편적으로 같다는 주장은 아니다.&lt;/p&gt;
&lt;/blockquote&gt;

&lt;h2&gt;
  
  
  Contract 4: Spark용 quality를 새로 만들지 않는다
&lt;/h2&gt;

&lt;p&gt;Spark용 quality suite를 별도로 구현하면 Python 경로와 조용히 달라질 수 있다. 그래서 S7은 Spark가 만든 silver/gold row를 driver로 collect한 뒤 기존 &lt;code&gt;build_quality_checks&lt;/code&gt;를 그대로 적용했다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;Spark transform
-&amp;gt; materialized silver/gold rows
-&amp;gt; existing checks
   - row count reconciliation
   - unit / defect conservation
   - numeric range
   - schema and business-date checks
-&amp;gt; pass일 때만 Iceberg write
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이 선택은 local bounded slice에는 적합하지만, distributed Spark-native quality evaluation은 아니다.&lt;/p&gt;

&lt;p&gt;quality가 실패하면 두 가지가 함께 성립해야 한다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;Iceberg snapshot을 만들지 않는다.
CLI가 non-zero로 종료되어 Airflow task도 실패한다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;처음 구현은 첫 번째만 지키고 &lt;code&gt;quality_failed&lt;/code&gt; JSON을 출력한 뒤 exit code &lt;code&gt;0&lt;/code&gt;으로 끝났다. 그러면 BashOperator는 데이터가 publish되지 않았는데도 task를 성공으로 볼 수 있다. review에서 이를 잡아 CLI가 &lt;code&gt;SystemExit(1)&lt;/code&gt;로 종료되도록 고쳤다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Contract 5: state 파일만 보고 skip하지 않는다
&lt;/h2&gt;

&lt;p&gt;Iceberg publish는 &lt;code&gt;overwritePartitions()&lt;/code&gt;를 쓴다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight python"&gt;&lt;code&gt;&lt;span class="nf"&gt;gold_dataframe&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;spark&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;gold_rows&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt; \
    &lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;writeTo&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;local.db.gold_daily_metrics&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt; \
    &lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;overwritePartitions&lt;/span&gt;&lt;span class="p"&gt;()&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;정정 source가 가진 &lt;code&gt;business_date&lt;/code&gt; partition만 교체하고, 다른 날짜 partition은 보존한다(같은 날짜 정정을 partition overwrite로 표현하는 이유는 B5에서 다뤘다). 같은 &lt;code&gt;table + business_date + source_hash&lt;/code&gt; 재실행은 새 snapshot을 만들지 않는다.&lt;/p&gt;

&lt;p&gt;여기에도 false skip이 있었다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;evidence state에는 이전 snapshot_id가 남아 있음
warehouse는 비워졌거나 새로 만들어짐
같은 source_hash가 다시 들어옴
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;state 파일의 hash만 보면 "이미 처리했다"고 skip하지만 실제 table은 비어 있을 수 있다. 그래서 skip 조건에 snapshot history membership을 추가했다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight python"&gt;&lt;code&gt;&lt;span class="n"&gt;same_source&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;previous_state&lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;source_hash&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;]&lt;/span&gt; &lt;span class="o"&gt;==&lt;/span&gt; &lt;span class="n"&gt;source_hash&lt;/span&gt;
&lt;span class="n"&gt;snapshot_exists&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="n"&gt;previous_state&lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;snapshot_id&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;]&lt;/span&gt; &lt;span class="ow"&gt;in&lt;/span&gt; &lt;span class="n"&gt;existing_snapshot_ids&lt;/span&gt;

&lt;span class="n"&gt;action&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;skip&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt; &lt;span class="k"&gt;if&lt;/span&gt; &lt;span class="n"&gt;same_source&lt;/span&gt; &lt;span class="ow"&gt;and&lt;/span&gt; &lt;span class="n"&gt;snapshot_exists&lt;/span&gt; &lt;span class="k"&gt;else&lt;/span&gt; &lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;write&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;snapshot expiry나 GC로 과거 snapshot이 history에서 사라지면 같은 partition을 한 번 더 쓰는 correct-but-extra rewrite가 생길 수 있다. 이번 local slice에서는 snapshot expiry를 실행하지 않았고, 데이터가 없는 상태에서 잘못 skip하는 것보다 재작성을 선택했다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Airflow에는 로직을 넣지 않았다
&lt;/h2&gt;

&lt;p&gt;Airflow DAG는 위 CLI 하나를 호출하는 single-task wrapper다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;Airflow
-&amp;gt; validated CLI command
-&amp;gt; adapter
-&amp;gt; Spark transform
-&amp;gt; quality gate
-&amp;gt; Iceberg publish
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;transform, quality, SparkSession, Iceberg write를 DAG body에 넣지 않았다. 로직은 일반 Python test와 Spark integration test로 검증하고, Airflow에서는 DAG import, command wiring, local &lt;code&gt;dags test&lt;/code&gt;만 확인했다.&lt;/p&gt;

&lt;p&gt;이는 production scheduler/executor나 운영 Airflow 검증이 아니다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Evidence
&lt;/h2&gt;

&lt;p&gt;독립 review 후 전체 검증 결과는 다음과 같다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;base Python environment:       90 passed, 14 skipped
Spark-visible environment:     99 passed, 5 skipped
S7 Spark integration:          14 passed
runtime state checks:          8/8 passed
isolated Airflow DagBag:       5 passed
local airflow dags test:       DagRun success, task exit 0
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;실제 local Iceberg 상태 전이도 확인했다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;다른 날짜 D2 publish
-&amp;gt; source A publish: snapshot 1 -&amp;gt; 2
-&amp;gt; 같은 source A retry: skipped, snapshot 그대로
-&amp;gt; correction source B: snapshot 2 -&amp;gt; 3
-&amp;gt; 대상 D1 units_produced=200으로 교체
-&amp;gt; D2 rows는 그대로 보존
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Spark gold의 &lt;code&gt;groupBy&lt;/code&gt; 실행 계획에서 &lt;code&gt;Exchange&lt;/code&gt;도 관찰했다. 이것은 shuffle이 발생한다는 local execution-plan 학습 evidence이지 성능이나 대규모 처리 성과는 아니다.&lt;/p&gt;

&lt;p&gt;구현과 테스트:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;&lt;a href="https://github.com/junhyun-dev/manufacturing-data-platform-mini/blob/main/src/manufacturing_data_platform/pipeline/spark_machine_event_batch.py" rel="noopener noreferrer"&gt;Spark machine-event batch code&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/junhyun-dev/manufacturing-data-platform-mini/blob/main/tests/test_spark_machine_event_batch.py" rel="noopener noreferrer"&gt;Engine parity and failure tests&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/junhyun-dev/manufacturing-data-platform-mini/blob/main/learn/reference-decisions/spark-engine-swap-contract.md" rel="noopener noreferrer"&gt;Engine-swap decision record&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/junhyun-dev/manufacturing-data-platform-mini/blob/main/VERIFICATION_LOG.md" rel="noopener noreferrer"&gt;Verification log&lt;/a&gt;&lt;/li&gt;
&lt;/ul&gt;

&lt;h2&gt;
  
  
  Limitations
&lt;/h2&gt;

&lt;p&gt;이번 글에서 증명하지 않은 것은 명확하다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;production / cluster Spark
대규모 성능·throughput 개선
full bronze/silver/gold Spark-Iceberg pipeline
Spark Structured Streaming 또는 direct Kafka-to-Iceberg
distributed Spark-native quality evaluation
concurrent Iceberg writer correctness
Iceberg commit과 JSON evidence write의 하나의 transaction
production Airflow operation
모든 float 입력에 대한 Python/Spark 반올림 동치
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;테이블도 local &lt;code&gt;gold_daily_metrics&lt;/code&gt; 하나뿐이다. 이 결과를 production lakehouse 구축이나 대규모 Spark 운영 경험으로 확대하지 않는다.&lt;/p&gt;

&lt;h2&gt;
  
  
  정리
&lt;/h2&gt;

&lt;p&gt;엔진 교체에서 먼저 비교해야 할 것은 코드 모양이 아니다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;같은 입력 identity인가?
같은 grain과 dedup 행을 선택하는가?
같은 metric을 만드는가?
같은 실패를 실패로 보고하는가?
같은 retry와 correction 상태 전이를 만드는가?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이번 slice에서 가장 값진 결과는 Spark 코드 자체보다 false-green 네 개를 계약과 테스트로 바꾼 일이었다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;반올림 parity 반례
stale snapshot state의 false skip
quality 실패의 exit 0
adapter provenance의 미지속
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;새 엔진을 붙였다는 말보다, 엔진이 바뀌어도 무엇이 같아야 하는지 설명하고 실패 반례로 증명하는 편이 더 강한 evidence가 된다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Sources
&lt;/h2&gt;

&lt;ul&gt;
&lt;li&gt;&lt;a href="https://spark.apache.org/docs/3.5.8/api/python/reference/pyspark.sql/api/pyspark.sql.functions.format_number.html" rel="noopener noreferrer"&gt;PySpark &lt;code&gt;format_number&lt;/code&gt;&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://spark.apache.org/docs/3.5.8/api/python/reference/pyspark.sql/api/pyspark.sql.DataFrameWriterV2.overwritePartitions.html" rel="noopener noreferrer"&gt;PySpark &lt;code&gt;DataFrameWriterV2.overwritePartitions&lt;/code&gt;&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://iceberg.apache.org/docs/latest/spark-writes/" rel="noopener noreferrer"&gt;Apache Iceberg Spark writes&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://airflow.apache.org/docs/apache-airflow/stable/howto/usage-cli.html" rel="noopener noreferrer"&gt;Apache Airflow CLI usage&lt;/a&gt;&lt;/li&gt;
&lt;/ul&gt;

</description>
      <category>dataengineering</category>
      <category>spark</category>
      <category>apacheiceberg</category>
      <category>airflow</category>
    </item>
    <item>
      <title>Kafka event를 바로 Spark로 보내지 않고 batch bridge를 둔 이유</title>
      <dc:creator>박준현</dc:creator>
      <pubDate>Sat, 18 Jul 2026 08:14:35 +0000</pubDate>
      <link>https://dev.to/junhyun-dev/kafka-eventreul-baro-sparkro-bonaeji-anhgo-batch-bridgereul-dun-iyu-3op6</link>
      <guid>https://dev.to/junhyun-dev/kafka-eventreul-baro-sparkro-bonaeji-anhgo-batch-bridgereul-dun-iyu-3op6</guid>
      <description>&lt;h1&gt;
  
  
  Kafka event를 바로 Spark로 보내지 않고 batch bridge를 둔 이유
&lt;/h1&gt;

&lt;p&gt;Kafka를 붙이고 나면 다음 선택은 자연스럽게 보인다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;Kafka
-&amp;gt; Spark Structured Streaming
-&amp;gt; Iceberg
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;하지만 이 프로젝트에서는 그렇게 하지 않았다.&lt;/p&gt;

&lt;p&gt;먼저 Kafka event를 payload와 &lt;code&gt;topic/partition/offset&lt;/code&gt;이 포함된 immutable JSONL로 보존했다. 그다음 한 &lt;code&gt;business_date&lt;/code&gt;의 accepted event만 결정적 CSV로 바꾸는 작은 bridge를 두고, 이미 검증한 batch quality/gold/Iceberg 경로를 재사용했다.&lt;/p&gt;

&lt;p&gt;이유는 간단하다.&lt;/p&gt;

&lt;blockquote&gt;
&lt;p&gt;이번 시나리오의 압력은 continuous window 처리가 아니라, event를 잃지 않고 replay 가능한 raw로 보존한 뒤 같은 accepted set이 trusted 결과를 두 배로 만들지 않게 하는 것이었다.&lt;/p&gt;
&lt;/blockquote&gt;

&lt;h2&gt;
  
  
  Scenario
&lt;/h2&gt;

&lt;p&gt;합성 설비 event가 local Kafka topic으로 들어온다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;producer
-&amp;gt; Kafka topic
-&amp;gt; bounded raw consumer
-&amp;gt; immutable JSONL landing
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;&lt;a href="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2Fzh5oe8us4k7kzqlkj6yh.png" class="article-body-image-wrapper"&gt;&lt;img src="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2Fzh5oe8us4k7kzqlkj6yh.png" alt="Kafka K1 and K1.5 runtime overview" width="800" height="500"&gt;&lt;/a&gt;&lt;/p&gt;

&lt;p&gt;consumer가 event를 읽었다는 사실과 downstream에 안전하게 보존했다는 사실은 같지 않다.&lt;/p&gt;

&lt;p&gt;가장 신경 쓴 실패 구간은 이것이었다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;1. Kafka record poll
2. raw landing write 완료
3. offset commit 전에 consumer crash
4. 같은 record 재전달
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;3번에서 죽으면 Kafka는 마지막 committed position부터 다시 전달한다. 이때 landing이 같은 record를 또 accepted로 쓰면 raw부터 중복된다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Decision Pressure: 무엇을 먼저 확정할 것인가
&lt;/h2&gt;

&lt;p&gt;Kafka offset commit과 local filesystem write는 하나의 transaction이 아니다.&lt;/p&gt;

&lt;p&gt;선택지는 두 가지였다.&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;순서&lt;/th&gt;
&lt;th&gt;장점&lt;/th&gt;
&lt;th&gt;실패 시 문제&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;offset commit -&amp;gt; file write&lt;/td&gt;
&lt;td&gt;중복 가능성이 작아 보임&lt;/td&gt;
&lt;td&gt;write 실패 시 record를 잃을 수 있음&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;file write -&amp;gt; offset commit&lt;/td&gt;
&lt;td&gt;record loss를 피함&lt;/td&gt;
&lt;td&gt;commit 전 crash 시 같은 record가 재전달됨&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;이 프로젝트는 두 번째를 선택했다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;poll bounded records
-&amp;gt; validate / classify
-&amp;gt; staging에 JSONL + manifest write
-&amp;gt; fsync
-&amp;gt; immutable batch directory로 rename
-&amp;gt; synchronous offset commit
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Kafka 공식 consumer 문서도 processing과 결합할 때 auto commit을 끄고 처리가 끝난 뒤 수동 commit하면, write 이후 commit 전 실패 구간에서 record가 재전달되는 at-least-once 특성이 생긴다고 설명한다.&lt;/p&gt;

&lt;p&gt;코드의 경계도 그대로 드러난다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight python"&gt;&lt;code&gt;&lt;span class="n"&gt;landing&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nf"&gt;land_records&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;
    &lt;span class="n"&gt;records&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;
    &lt;span class="n"&gt;output_dir&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;
    &lt;span class="n"&gt;simulate_crash_after_rename&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="n"&gt;simulate_crash_after_landing&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;
&lt;span class="p"&gt;)&lt;/span&gt;

&lt;span class="n"&gt;offsets&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="p"&gt;[&lt;/span&gt;
    &lt;span class="nc"&gt;TopicPartition&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;item&lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;topic&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;],&lt;/span&gt; &lt;span class="n"&gt;item&lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;partition&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;],&lt;/span&gt; &lt;span class="n"&gt;item&lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;next_offset&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;])&lt;/span&gt;
    &lt;span class="k"&gt;for&lt;/span&gt; &lt;span class="n"&gt;item&lt;/span&gt; &lt;span class="ow"&gt;in&lt;/span&gt; &lt;span class="n"&gt;landing&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;committable_offsets&lt;/span&gt;
&lt;span class="p"&gt;]&lt;/span&gt;
&lt;span class="n"&gt;consumer&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;commit&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;offsets&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="n"&gt;offsets&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt; &lt;span class="n"&gt;asynchronous&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="bp"&gt;False&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;핵심은 재전달 자체를 없애는 것이 아니다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;at-least-once 재전달은 허용한다.
이미 durable한 (topic, partition, offset)은 다시 accepted하지 않는다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h2&gt;
  
  
  Failure Injection: 정말 복구되는가
&lt;/h2&gt;

&lt;p&gt;정상 경로만 실행해서는 이 계약을 증명할 수 없다. 그래서 offset &lt;code&gt;3&lt;/code&gt;을 landing한 직후, commit 전에 의도적으로 예외를 발생시켰다.&lt;/p&gt;

&lt;p&gt;&lt;a href="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2Fhgffri89aqpxyxx1uij8.png" class="article-body-image-wrapper"&gt;&lt;img src="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2Fhgffri89aqpxyxx1uij8.png" alt="Kafka landing-before-commit failure recovery" width="800" height="500"&gt;&lt;/a&gt;&lt;/p&gt;

&lt;p&gt;실제 검증 결과:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;produced coordinates: 5
accepted events: 4
quarantined events: 1
persisted coordinates: 5

injected crash: observed
redelivered coordinate: 1
recovery landing status: reused
accepted total after recovery: 4
committed next offset: 4
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;재전달은 일어났지만 accepted set은 늘지 않았다.&lt;/p&gt;

&lt;p&gt;별도의 bounded replay도 실행했다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;replay offsets: 0..3
reused coordinates: 4
normal consumer-group commit: not performed
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;replay가 정상 consumer-group progress를 바꾸지 않는 것도 별도 계약으로 확인했다.&lt;/p&gt;

&lt;h2&gt;
  
  
  왜 바로 Spark Structured Streaming으로 가지 않았나
&lt;/h2&gt;

&lt;p&gt;K1 raw landing이 끝나도 기존 pipeline과는 입력 형식이 다르다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;K1 output: accepted JSONL + manifest
기존 pipeline input: one business_date CSV
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;여기서 세 가지 선택지를 비교했다.&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Option&lt;/th&gt;
&lt;th&gt;장점&lt;/th&gt;
&lt;th&gt;이번 slice의 비용&lt;/th&gt;
&lt;th&gt;판단&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;Spark Structured Streaming&lt;/td&gt;
&lt;td&gt;window, checkpoint, continuous 처리 가능&lt;/td&gt;
&lt;td&gt;아직 없는 latency/window 요구까지 운영해야 함&lt;/td&gt;
&lt;td&gt;Backlog&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;direct Kafka -&amp;gt; Iceberg sink&lt;/td&gt;
&lt;td&gt;경로가 짧음&lt;/td&gt;
&lt;td&gt;기존 quality/catalog gate를 우회&lt;/td&gt;
&lt;td&gt;제외&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;deterministic batch bridge&lt;/td&gt;
&lt;td&gt;기존 검증 자산 재사용&lt;/td&gt;
&lt;td&gt;adapter contract가 하나 필요&lt;/td&gt;
&lt;td&gt;선택&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;K1.5 bridge는 한 날짜의 accepted event를 &lt;code&gt;(topic, partition, offset)&lt;/code&gt; 순서로 정렬하고, 고정 header와 newline으로 canonical CSV를 만든다.&lt;/p&gt;

&lt;p&gt;여기서 결정적이라는 말은 &lt;strong&gt;주어진 하나의 immutable landing에서 같은 accepted set을 다시 선택했을 때&lt;/strong&gt; 같은 bytes가 나온다는 뜻이다. producer를 새로 실행해 Kafka timestamp나 record fingerprint가 바뀌면 provenance가 달라지므로 새 &lt;code&gt;source_hash&lt;/code&gt;가 만들어진다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight python"&gt;&lt;code&gt;&lt;span class="n"&gt;selected&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;sort&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;key&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;&lt;span class="k"&gt;lambda&lt;/span&gt; &lt;span class="n"&gt;item&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt; &lt;span class="n"&gt;item&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="n"&gt;sort_key&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
&lt;span class="n"&gt;csv_bytes&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nf"&gt;canonical_csv_bytes&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;selected&lt;/span&gt;&lt;span class="p"&gt;)&lt;/span&gt;
&lt;span class="n"&gt;source_hash&lt;/span&gt; &lt;span class="o"&gt;=&lt;/span&gt; &lt;span class="nf"&gt;sha256&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="n"&gt;csv_bytes&lt;/span&gt;&lt;span class="p"&gt;).&lt;/span&gt;&lt;span class="nf"&gt;hexdigest&lt;/span&gt;&lt;span class="p"&gt;()&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이 CSV에는 business column뿐 아니라 다음 provenance도 들어간다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;event_id
schema_version
kafka_topic
kafka_partition
kafka_offset
kafka_key
kafka_timestamp_ms
kafka_record_fingerprint
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;그리고 CSV의 SHA-256을 기존 pipeline의 &lt;code&gt;source_hash&lt;/code&gt;로 그대로 사용한다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;same accepted Kafka set from one immutable landing
-&amp;gt; same canonical CSV bytes
-&amp;gt; same source_hash
-&amp;gt; same adapter version reused
-&amp;gt; existing lakehouse run skipped
-&amp;gt; Iceberg retry creates no new snapshot
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;새 idempotency 체계를 하나 더 만든 것이 아니라, Kafka provenance를 기존 source identity 안으로 연결했다.&lt;/p&gt;

&lt;p&gt;이 경로에는 서로 다른 다섯 개의 identity가 등장한다. 헷갈리면 안 되는 부분이다.&lt;/p&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Identity&lt;/th&gt;
&lt;th&gt;무엇을 가리키나&lt;/th&gt;
&lt;th&gt;누가 정하나&lt;/th&gt;
&lt;th&gt;같은 accepted set 재실행 시&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;&lt;code&gt;(topic, partition, offset)&lt;/code&gt;&lt;/td&gt;
&lt;td&gt;Kafka log 안의 transport 위치&lt;/td&gt;
&lt;td&gt;Kafka broker&lt;/td&gt;
&lt;td&gt;coordinate 그대로 재사용&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;&lt;code&gt;event_id&lt;/code&gt;&lt;/td&gt;
&lt;td&gt;business event 한 건의 정체&lt;/td&gt;
&lt;td&gt;producer가 생성&lt;/td&gt;
&lt;td&gt;같은 event_id는 accepted set을 늘리지 않음&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;adapter &lt;code&gt;source_hash&lt;/code&gt;
&lt;/td&gt;
&lt;td&gt;선택된 한 날짜 batch 입력의 정체&lt;/td&gt;
&lt;td&gt;canonical CSV의 SHA-256&lt;/td&gt;
&lt;td&gt;같은 값 → adapter version 재사용&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;pipeline &lt;code&gt;run_id&lt;/code&gt;
&lt;/td&gt;
&lt;td&gt;lakehouse 실행 한 번의 정체&lt;/td&gt;
&lt;td&gt;pipeline이 실행 시 생성&lt;/td&gt;
&lt;td&gt;재실행은 기존 run_id를 skip으로 재사용&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;Iceberg &lt;code&gt;snapshot_id&lt;/code&gt;
&lt;/td&gt;
&lt;td&gt;gold table commit 한 번의 정체&lt;/td&gt;
&lt;td&gt;Iceberg가 write 시 생성&lt;/td&gt;
&lt;td&gt;새 snapshot을 만들지 않음&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;핵심은 transport 위치(offset)와 business 정체(event_id)와 batch 정체(source_hash)를 분리해 두었기 때문에, 재전달이나 재실행이 어느 층에서 일어나도 그 위쪽 trusted 결과가 늘지 않는다는 것이다.&lt;/p&gt;

&lt;h2&gt;
  
  
  End-to-End Evidence
&lt;/h2&gt;

&lt;p&gt;&lt;a href="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2Fle3cs89hn300ctuwwmd1.png" class="article-body-image-wrapper"&gt;&lt;img src="https://media2.dev.to/dynamic/image/width=800%2Cheight=%2Cfit=scale-down%2Cgravity=auto%2Cformat=auto/https%3A%2F%2Fdev-to-uploads.s3.us-east-2.amazonaws.com%2Fuploads%2Farticles%2Fle3cs89hn300ctuwwmd1.png" alt="Kafka landing to batch and Iceberg rerun evidence" width="800" height="500"&gt;&lt;/a&gt;&lt;/p&gt;

&lt;p&gt;2026-07-16 local runtime 결과:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;Kafka: 4.3.1, one broker / one topic / one partition
K1 broker verification: passed
K1.5 checks: 11 passed

selected accepted events: 4
adapter: created -&amp;gt; reused
lakehouse: processed -&amp;gt; skipped
quality checks: 8 / 8 pass
silver rows: 4
gold rows: 1
gold totals: units=100, defects=6

Iceberg: published -&amp;gt; skipped
snapshot count: 1 -&amp;gt; 1
snapshot ID unchanged: true
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;또한 adapter는 요청 날짜에 accepted event가 하나도 없거나 manifest와 JSONL이 어긋나면 기존 pipeline을 호출하기 전에 실패한다. 그렇지 않으면 기존 pipeline의 synthetic sample 생성 경로가 실제 입력이 없는데도 성공처럼 보이는 false green을 만들 수 있기 때문이다.&lt;/p&gt;

&lt;h2&gt;
  
  
  What This Changed in the Design
&lt;/h2&gt;

&lt;p&gt;처음에는 Kafka 다음에 Spark Structured Streaming을 붙이는 순서를 생각했다.&lt;/p&gt;

&lt;p&gt;실제 시나리오와 failure contract를 따라가자 순서가 바뀌었다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;K1   Kafka -&amp;gt; immutable raw landing
K1.5 raw landing -&amp;gt; deterministic batch adapter -&amp;gt; existing quality/gold/Iceberg
K2   Spark Structured Streaming은 window/latency pressure가 생길 때만
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;도구의 자연스러운 나열보다 현재 문제를 닫는 가장 작은 vertical slice를 선택한 것이다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Limitations
&lt;/h2&gt;

&lt;p&gt;이 글이 증명하는 범위:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;local Kafka one broker / one topic / one partition
bounded producer and consumer
landing-before-commit failure injection
coordinate-based raw reuse
bounded replay without normal group commit
one business_date deterministic batch bridge
existing local quality/gold/Iceberg retry contract reuse
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;증명하지 않은 범위:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;continuous streaming service
multi-partition ordering / rebalance
multi-broker HA
Spark Structured Streaming
direct Kafka-to-Iceberg sink
end-to-end exactly-once
power-loss durability
production security / scale / operations
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;특히 &lt;code&gt;os.replace&lt;/code&gt;와 &lt;code&gt;fsync&lt;/code&gt; 순서는 local Linux filesystem에서만 검증했다. in-process 예외를 주입한 것이므로 전원 손실까지 견디는 crash consistency를 주장하지 않는다.&lt;/p&gt;

&lt;h2&gt;
  
  
  정리
&lt;/h2&gt;

&lt;p&gt;Kafka를 썼다는 사실보다 중요한 것은 &lt;strong&gt;언제 offset을 확정하고, 재전달을 어떤 identity로 흡수하며, 기존 trusted pipeline과 어떻게 연결하는가&lt;/strong&gt;였다.&lt;/p&gt;

&lt;p&gt;이번 구현의 핵심은 세 줄이다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;durable landing 전에 offset을 commit하지 않는다.
재전달된 Kafka coordinate는 accepted set을 늘리지 않는다.
같은 accepted set은 기존 source_hash 계약을 통해 gold와 snapshot도 늘리지 않는다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;코드와 전체 walkthrough: &lt;a href="https://github.com/junhyun-dev/manufacturing-data-platform-mini/tree/main/docs/portfolio/kafka-k1-k1-5" rel="noopener noreferrer"&gt;manufacturing-data-platform-mini&lt;/a&gt;&lt;/p&gt;

&lt;h2&gt;
  
  
  Sources
&lt;/h2&gt;

&lt;ul&gt;
&lt;li&gt;&lt;a href="https://kafka.apache.org/43/javadoc/org/apache/kafka/clients/consumer/KafkaConsumer.html" rel="noopener noreferrer"&gt;Apache Kafka 4.3.1 KafkaConsumer: offsets, manual commits, at-least-once delivery&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://docs.python.org/3/library/os.html#os.replace" rel="noopener noreferrer"&gt;Python &lt;code&gt;os.replace&lt;/code&gt;&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://docs.python.org/3/library/os.html#os.fsync" rel="noopener noreferrer"&gt;Python &lt;code&gt;os.fsync&lt;/code&gt;&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://spark.apache.org/docs/latest/streaming/index.html" rel="noopener noreferrer"&gt;Spark Structured Streaming guide&lt;/a&gt;&lt;/li&gt;
&lt;/ul&gt;

</description>
      <category>dataengineering</category>
      <category>kafka</category>
      <category>python</category>
      <category>apacheiceberg</category>
    </item>
    <item>
      <title>AI Memory OSS를 프롬프트가 아니라 백엔드 시스템으로 읽는 10가지 질문</title>
      <dc:creator>박준현</dc:creator>
      <pubDate>Fri, 17 Jul 2026 05:26:15 +0000</pubDate>
      <link>https://dev.to/junhyun-dev/ai-memory-ossreul-peurompeuteuga-anira-baegendeu-siseutemeuro-ilgneun-10gaji-jilmun-1h9h</link>
      <guid>https://dev.to/junhyun-dev/ai-memory-ossreul-peurompeuteuga-anira-baegendeu-siseutemeuro-ilgneun-10gaji-jilmun-1h9h</guid>
      <description>&lt;h1&gt;
  
  
  AI Memory OSS를 프롬프트가 아니라 백엔드 시스템으로 읽는 10가지 질문
&lt;/h1&gt;

&lt;p&gt;AI memory를 처음 보면 prompt, embedding, vector search부터 눈에 들어온다. 하지만 기능을 운영 가능한 서비스로 만들려면 그 앞뒤의 상태와 실패를 함께 봐야 한다. 원문은 언제 저장되는지, 느린 추론은 어디서 실행되는지, 파생된 기억은 무엇을 근거로 하는지, 검색 인덱스와 DB가 어긋나면 어떻게 복구하는지까지 답해야 한다.&lt;/p&gt;

&lt;p&gt;이 글은 AI memory OSS인 Honcho를 pinned source와 변경 이력으로 읽으면서 만든 질문 지도다. Honcho의 모든 설계를 설명하거나 평가하려는 글이 아니라, 기능 설명을 실제 backend contract로 내려가기 위해 무엇을 물어야 하는지 정리한다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Scenario
&lt;/h2&gt;

&lt;p&gt;사용자와 Agent가 오랫동안 대화하는 서비스를 생각해 본다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;사용자가 메시지를 보낸다.
-&amp;gt; 원문이 저장된다.
-&amp;gt; 느린 LLM 작업이 기억 후보를 만든다.
-&amp;gt; Agent가 나중에 필요한 기억을 검색한다.
-&amp;gt; 사용자는 말을 바꾸거나 삭제를 요청한다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;짧은 demo라면 최근 대화와 vector search만으로도 동작할 수 있다. 하지만 서비스가 오래 실행되면 질문이 달라진다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;LLM 호출 중 API와 DB resource는 어떻게 되는가?
원문과 LLM이 만든 주장은 같은 데이터인가?
worker가 죽으면 어떤 상태부터 다시 시작하는가?
같은 기억이 중복 생성되면 무엇을 기준으로 막는가?
DB에는 있는데 vector index에는 없는 기억을 누가 복구하는가?
삭제는 DB row만 지우면 끝나는가?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;그래서 AI memory를 다음과 같은 lifecycle로 읽었다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;write                                      (Q1)
-&amp;gt; durable raw state                       (Q2)
-&amp;gt; queue / worker                          (Q3)
-&amp;gt; derived claim + provenance              (Q2, Q4)
   provenance = 원본까지의 근거 연결
-&amp;gt; vector index sync                       (Q7)
-&amp;gt; scoped retrieval                        (Q5, Q6)
-&amp;gt; correction / deletion                   (Q8)
-&amp;gt; observability / cost control            (Q9)

schema evolution은 이 lifecycle 전체를 가로지른다. (Q10)
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;각 화살표는 구현 단계인 동시에 실패할 수 있는 경계다.&lt;/p&gt;

&lt;h2&gt;
  
  
  1. API boundary: 요청 안에서 어디까지 끝낼 것인가?
&lt;/h2&gt;



&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;HTTP 응답 전에 반드시 저장해야 하는 상태는 무엇인가?
느린 embedding/LLM 호출은 어디서 끊어야 하는가?
외부 호출을 기다리는 동안 DB transaction이나 connection을 점유하는가?
요청 실패와 background 작업 실패를 어떻게 구분하는가?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;비동기 API를 사용한다고 resource가 자동으로 짧게 유지되는 것은 아니다. Honcho의 한 변경을 추적했을 때도 핵심은 &lt;code&gt;async&lt;/code&gt; 문법이 아니라 외부 호출과 DB session/transaction scope를 분리한 것이었다. 이 사례는 &lt;a href="https://dev.to/junhyun-dev/neurin-llm-hocul-jung-db-connectioneul-jabji-anhneun-iyu-3abg"&gt;첫 번째 글&lt;/a&gt;에서 자세히 다뤘다.&lt;/p&gt;

&lt;h2&gt;
  
  
  2. State model: 원문과 파생된 기억의 grain(한 row가 의미하는 단위)은 무엇인가?
&lt;/h2&gt;



&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;message와 memory document는 왜 분리되는가?
어떤 row가 원본이고 어떤 row가 LLM이 만든 claim인가?
claim은 어떤 source를 가리키는가?
사실, 추론, 패턴, 모순을 같은 상태로 취급하는가?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;AI가 만든 문장을 원문과 같은 진실로 저장하면 나중에 수정, 재계산, 설명이 어려워진다. Honcho의 pinned model도 &lt;code&gt;Message&lt;/code&gt;와 &lt;code&gt;Document&lt;/code&gt;를 분리하고, document에 &lt;code&gt;level&lt;/code&gt;과 &lt;code&gt;source_ids&lt;/code&gt;를 둔다. 여기서 가져갈 일반적인 질문은 특정 테이블 이름이 아니라 다음 계약이다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;derived claim은 raw record와 다른 identity와 provenance를 가져야 하는가?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h2&gt;
  
  
  3. Queue/worker: queue item보다 논리적 작업의 identity가 있는가?
&lt;/h2&gt;



&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;retry해도 같은 일임을 어떻게 식별하는가?
같은 사용자나 session의 작업은 어디까지 직렬화하는가?
다른 작업은 병렬로 실행할 수 있는가?
처리 중 worker가 죽으면 어떤 claim과 checkpoint가 남는가?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Queue를 도입했다는 사실만으로 중복 처리와 순서 문제가 해결되지는 않는다. payload보다 먼저 &lt;code&gt;이 시스템에서 같은 일은 무엇인가?&lt;/code&gt;를 정의해야 한다. Honcho의 queue model에는 &lt;code&gt;work_unit_key&lt;/code&gt;(같은 논리적 작업임을 식별하는 key)가 있고, 일부 pending 작업에는 이 key를 기준으로 한 unique constraint가 있다. 다만 enqueue부터 recovery까지의 전체 계약은 후속 source audit 대상으로 남겨뒀다.&lt;/p&gt;

&lt;h2&gt;
  
  
  4. LLM boundary: 모델이 판단하는 것과 코드가 강제하는 것은 무엇인가?
&lt;/h2&gt;



&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;LLM은 어떤 판단만 맡는가?
출력이 비어 있거나 구조가 깨지면 어떻게 하는가?
source 없는 claim을 허용하는가?
성공은 API 응답 수신인가, schema 통과인가, 저장 완료인가?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;LLM에게 유연한 판단을 맡기더라도 identity, 허용 상태, provenance, checkpoint 같은 계약은 결정적으로 검증할 필요가 있다. 모델의 좋은 답을 기대하는 것과 시스템이 잘못된 상태를 거부하는 것은 다른 책임이다.&lt;/p&gt;

&lt;h2&gt;
  
  
  5. Identity/authz: actor, tenant, resource, subject가 같은가?
&lt;/h2&gt;

&lt;p&gt;여기서 tenant는 고객 또는 workspace 단위의 격리 경계를 뜻한다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;요청을 보낸 actor는 누구인가?
어느 workspace의 resource에 접근하는가?
누가 누구에 대해 만든 기억인가?
DB query와 vector search에 같은 tenant scope가 적용되는가?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Agent와 사용자를 하나의 identity model로 표현하면 도메인 모델은 유연해질 수 있다. 반면 &lt;code&gt;요청자&lt;/code&gt;, &lt;code&gt;기억의 관찰자&lt;/code&gt;, &lt;code&gt;기억의 대상&lt;/code&gt;, &lt;code&gt;resource owner&lt;/code&gt;가 서로 다른 축이 될 수 있다. 편리한 도메인 모델이 권한 검증까지 자동으로 단순화한다고 가정하면 안 된다.&lt;/p&gt;

&lt;h2&gt;
  
  
  6. Retrieval: 저장된 기억 중 무엇이 실제 답변 후보가 되는가?
&lt;/h2&gt;



&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;semantic search 전에 어떤 tenant/session/level filter를 적용하는가?
원문 기반 claim과 상위 추론을 같은 우선순위로 검색하는가?
deleted, stale, sync-pending 상태는 후보에서 제외되는가?
응답에 source와 reasoning path를 함께 전달할 수 있는가?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;기억을 만드는 pipeline과 Agent가 사용하는 serving path는 별개의 계약이다. 저장 품질이 좋아도 retrieval filter가 잘못되면 다른 tenant의 기억, 삭제된 기억, 근거가 약한 추론이 답변에 들어갈 수 있다. Honcho의 retrieval filter 구성은 아직 별도 source audit 전이다.&lt;/p&gt;

&lt;h2&gt;
  
  
  7. Index consistency: DB와 vector store의 반쪽 성공을 어떻게 복구하는가?
&lt;/h2&gt;



&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;DB commit 후 vector upsert가 실패하면 무엇이 남는가?
DB에는 있지만 index에는 없는 row를 어떻게 찾는가?
재시도 횟수와 마지막 성공 시점을 기록하는가?
vector backend를 바꿔도 serving contract가 유지되는가?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Honcho의 pinned source에는 document의 &lt;code&gt;sync_state&lt;/code&gt;, &lt;code&gt;last_sync_at&lt;/code&gt;, &lt;code&gt;sync_attempts&lt;/code&gt;와 DB와 vector index 사이의 어긋난 상태를 맞추는 background 작업인 reconciler가 존재한다. 이 작업은 미동기 항목의 재시도와 삭제 잔여물 정리를 담당한다. 이것은 중요한 조사 단서지만, 모든 crash window와 복구 의미를 확인했다는 뜻은 아니다. DB-vector consistency는 별도의 teardown으로 검증할 예정이다.&lt;/p&gt;

&lt;h2&gt;
  
  
  8. Deletion/retention: 언제 삭제가 완료됐다고 말할 수 있는가?
&lt;/h2&gt;



&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;message를 삭제하면 파생된 document도 함께 처리되는가?
DB, vector index, cache, backup 중 어디까지 삭제해야 하는가?
soft delete와 hard delete 사이의 상태는 누가 소유하는가?
외부 index 삭제가 실패하면 DB row를 보존하는가?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;AI memory에서 삭제는 개인정보 문제이면서 분산 상태 정합성 문제다. Honcho source에는 &lt;code&gt;deleted_at&lt;/code&gt;으로 표시한 뒤 reconciler가 외부 vector 삭제와 DB hard delete를 처리하는 경로가 있다. 이 경로 역시 lifecycle 전체를 별도 audit하기 전에는 설계 완료 여부를 단정하지 않는다.&lt;/p&gt;

&lt;h2&gt;
  
  
  9. Observability/cost: 운영자는 어떤 질문에 답할 수 있는가?
&lt;/h2&gt;



&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;queue backlog와 stale work를 볼 수 있는가?
workspace별 memory 생성 지연을 알 수 있는가?
LLM latency, token usage, retry 비용을 추적하는가?
특정 memory가 왜 생성됐는지 source까지 추적할 수 있는가?
DB connection pool과 외부 provider 장애를 구분할 수 있는가?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;관측성은 dashboard 제품명을 고르는 문제가 아니다. 먼저 운영자가 복구와 비용 판단을 위해 답해야 할 질문을 정하고, 그 질문에 필요한 durable state와 event를 남겨야 한다. Honcho의 pinned source에도 message별 &lt;code&gt;token_count&lt;/code&gt;와 sync 재시도 횟수 같은 단서 상태가 있지만, 운영 관측 전체는 아직 검증하지 않았다.&lt;/p&gt;

&lt;h2&gt;
  
  
  10. Schema evolution: 처음 정한 memory model이 틀리면 어떻게 바꾸는가?
&lt;/h2&gt;



&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;새 identity나 document level을 기존 데이터에 어떻게 적용하는가?
old API/worker와 new schema가 동시에 존재할 수 있는가?
vector namespace나 ID 규칙 변경을 어떻게 backfill하는가?
rollback하면 새로 생성된 data shape을 어떻게 읽는가?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Agent 제품의 개념은 계속 바뀔 수 있다. 따라서 모델을 잘 정하는 능력만큼, 기존 데이터를 잃지 않고 계약을 바꾸는 migration 경로도 중요하다. Honcho에는 migration 이력이 존재하지만 이 글에서는 그 경로를 검증하지 않았다.&lt;/p&gt;

&lt;h2&gt;
  
  
  질문을 실제 조사로 내리는 방법
&lt;/h2&gt;

&lt;p&gt;이 질문들을 모든 기능에 한꺼번에 적용하면 거대한 체크리스트만 남는다. 이 글에서는 다음 순서로 범위를 줄였다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;1. 지금 보는 기능의 actor를 정한다.
2. durable state transition을 시간순으로 적는다.
3. 그 전이에서 실제로 발생하는 pressure 하나를 고른다.
4. 코드, migration, test, commit에서 현재 답을 찾는다.
5. 내 시스템에는 copy / simplify / avoid 중 무엇을 할지 결정한다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;예를 들어 dream scheduling을 볼 때는 &lt;code&gt;AI가 기억을 만든다&lt;/code&gt;에서 멈추지 않는다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;actor: scheduler / worker
input: source-like explicit documents
in-flight owner: pending queue item
completed progress: success checkpoint
pressure: 비용, retry, source/derived feedback loop
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이렇게 상태의 owner를 나누자, 파생된 출력이 다음 실행의 eligibility를 키우는 문제와 enqueue 시점에 checkpoint를 먼저 전진시키는 문제가 보였다. Honcho가 이를 어떻게 변경했는지는 &lt;a href="https://dev.to/junhyun-dev/ai-memoryga-jagi-culryeogeul-dasi-ibryeogeuro-seji-anhge-haneun-beob-2k0g"&gt;두 번째 글&lt;/a&gt;에서 추적했다.&lt;/p&gt;

&lt;h2&gt;
  
  
  이 지도가 증명하지 않는 것
&lt;/h2&gt;

&lt;ul&gt;
&lt;li&gt;이 글은 Honcho의 공개 코드와 변경 이력을 읽으며 만든 분석 틀이다. Honcho를 직접 구현, 운영하거나 기여한 경험이 아니다.&lt;/li&gt;
&lt;li&gt;분석 기준은 &lt;code&gt;85239a69b262c944de3c35900b91c88ba9b84f1a&lt;/code&gt;로 고정했다.&lt;/li&gt;
&lt;li&gt;열 가지 질문이 모든 AI memory 서비스에 동일한 우선순위로 필요하다고 주장하지 않는다.&lt;/li&gt;
&lt;li&gt;DB-vector consistency, deletion, observability, migration은 source에서 조사 단서를 확인했지만, 아직 lifecycle 전체와 runtime behavior를 검증하지 않았다.&lt;/li&gt;
&lt;li&gt;이 질문 지도는 security review, privacy review, production load test를 대체하지 않는다.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;다음 teardown에서는 현재 source evidence가 충분한 주제부터 multi-tenant authz와 work-unit serialization을 검토한다. Retrieval, index consistency, deletion은 실제 실패 경계와 테스트를 확인한 뒤에만 공개 글로 승격할 예정이다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Sources
&lt;/h2&gt;

&lt;ul&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho" rel="noopener noreferrer"&gt;Honcho repository&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/tree/85239a69b262c944de3c35900b91c88ba9b84f1a" rel="noopener noreferrer"&gt;Pinned Honcho source tree&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/src/models.py#L206-L528" rel="noopener noreferrer"&gt;Pinned message, document, queue state models&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/src/routers/messages.py#L88-L141" rel="noopener noreferrer"&gt;Pinned message write and enqueue boundary&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/src/crud/document.py#L197-L421" rel="noopener noreferrer"&gt;Pinned document retrieval paths&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/src/reconciler/sync_vectors.py#L78-L455" rel="noopener noreferrer"&gt;Pinned vector sync reconciler&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/src/crud/document.py#L1031-L1106" rel="noopener noreferrer"&gt;Pinned soft-delete cleanup path&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://dev.to/junhyun-dev/neurin-llm-hocul-jung-db-connectioneul-jabji-anhneun-iyu-3abg"&gt;느린 LLM 호출 중 DB connection을 잡지 않는 이유&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://dev.to/junhyun-dev/ai-memoryga-jagi-culryeogeul-dasi-ibryeogeuro-seji-anhge-haneun-beob-2k0g"&gt;AI memory가 자기 출력을 다시 입력으로 세지 않게 하는 법&lt;/a&gt;&lt;/li&gt;
&lt;/ul&gt;

</description>
      <category>ai</category>
      <category>backend</category>
      <category>architecture</category>
      <category>opensource</category>
    </item>
    <item>
      <title>AI memory가 자기 출력을 다시 입력으로 세지 않게 하는 법</title>
      <dc:creator>박준현</dc:creator>
      <pubDate>Thu, 16 Jul 2026 13:55:57 +0000</pubDate>
      <link>https://dev.to/junhyun-dev/ai-memoryga-jagi-culryeogeul-dasi-ibryeogeuro-seji-anhge-haneun-beob-2k0g</link>
      <guid>https://dev.to/junhyun-dev/ai-memoryga-jagi-culryeogeul-dasi-ibryeogeuro-seji-anhge-haneun-beob-2k0g</guid>
      <description>&lt;h1&gt;
  
  
  AI memory가 자기 출력을 다시 입력으로 세지 않게 하는 법
&lt;/h1&gt;

&lt;p&gt;AI memory 시스템은 대화에서 기억을 만들고, 그 기억을 다시 요약해 더 높은 수준의 결론을 만들 수 있다. 여기서 원본에 가까운 입력과 시스템이 만든 파생 결과를 구분하지 않으면, 파생 결과가 다음 작업의 trigger를 키우는 self-reinforcing loop가 생길 수 있다.&lt;/p&gt;

&lt;p&gt;이 글은 AI memory OSS Honcho의 PR #573과 pinned source를 읽으며, 이런 feedback pressure를 입력 eligibility, 성공 checkpoint, in-flight state로 나눠 해결한 방식을 추적한 기록이다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Scenario
&lt;/h2&gt;

&lt;p&gt;Honcho의 dreamer는 collection(observer-observed 쌍마다 memory document가 쌓이는 묶음)에 있는 document를 바탕으로 더 높은 수준의 관찰을 만든다. Document에는 여러 level이 있다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;explicit       원문에서 직접 도출된 관찰
deductive      기존 관찰에서 연역한 결과
inductive      여러 관찰에서 귀납한 결과
contradiction  서로 충돌하는 관찰을 표현한 결과
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;dreamer가 실행되면 deductive, inductive, contradiction 같은 새로운 document가 생길 수 있다. 동시에 scheduler는 collection에 document가 일정 개수 이상 추가되면 다음 dream을 예약한다.&lt;/p&gt;

&lt;p&gt;문제는 &lt;strong&gt;무엇을 "새 입력"으로 셀 것인가&lt;/strong&gt;다.&lt;/p&gt;

&lt;h2&gt;
  
  
  변경 전: 모든 document가 같은 trigger에 들어갔다
&lt;/h2&gt;

&lt;p&gt;PR #573 직전의 &lt;code&gt;check_and_schedule_dream&lt;/code&gt;은 workspace와 observer/observed만으로 document 수를 계산했다. &lt;code&gt;level&lt;/code&gt; 조건이 없었기 때문에 explicit input뿐 아니라 이전 dream이 만든 derived output도 threshold에 포함됐다.&lt;/p&gt;

&lt;p&gt;예를 들어 threshold가 50일 때 다음 collection을 생각할 수 있다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;explicit 30
deductive 40
inductive 10
----------------
전체 document 80
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;새로운 explicit 관찰은 30개뿐이지만 전체 count는 80이므로 scheduler 조건을 통과할 수 있다. 이전 dream의 출력이 다음 dream의 trigger를 키우는 구조다.&lt;/p&gt;

&lt;p&gt;이것만으로 production에서 무한 실행이 발생했다고 단정할 수는 없다. 시간 간격 제한과 queue deduplication 같은 다른 gate도 있기 때문이다. 여기서 source와 regression test로 확인한 것은 &lt;strong&gt;derived output이 trigger count를 스스로 증가시킬 수 있었던 eligibility 오류&lt;/strong&gt;다.&lt;/p&gt;

&lt;h2&gt;
  
  
  첫 번째 경계: source-like input만 센다
&lt;/h2&gt;

&lt;p&gt;변경 후 threshold query는 &lt;code&gt;Document.level == "explicit"&lt;/code&gt; 조건을 사용한다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;current_explicit_count - last_dream_document_count &amp;gt;= threshold
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;deductive, inductive, contradiction document는 검색과 추론에는 계속 사용할 수 있지만, 다음 consolidation을 시작시키는 새 입력으로는 세지 않는다. 저장 여부와 scheduling eligibility를 분리한 것이다.&lt;/p&gt;

&lt;p&gt;pinned current code에서는 실행할 session을 고를 때도 최신 explicit document만 본다. threshold는 explicit 기준 집합으로 계산하면서 실행 context는 더 최신인 derived document의 session에서 가져오는 비대칭을 피하기 위한 조치다.&lt;/p&gt;

&lt;h2&gt;
  
  
  두 번째 경계: enqueue가 아니라 성공 뒤에 checkpoint한다
&lt;/h2&gt;

&lt;p&gt;변경 전 &lt;code&gt;enqueue_dream&lt;/code&gt;은 queue에 작업을 넣으면서 다음 두 값을 바로 기록했다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;last_dream_document_count
last_dream_at
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;하지만 enqueue는 작업이 실행되거나 유효한 결과를 만들었다는 뜻이 아니다. 실패한 작업이 baseline을 먼저 전진시키면 같은 corpus를 다시 시도할 근거를 잃을 수 있다.&lt;/p&gt;

&lt;p&gt;현재 코드는 enqueue 단계에서 dream metadata를 건드리지 않는다. &lt;code&gt;process_dream&lt;/code&gt;이 non-null result를 받은 뒤에만 row lock을 잡고 현재 explicit count를 다시 계산해, count와 timestamp를 함께 기록한다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;enqueue          in-flight 상태만 생성
run_dream 실패   checkpoint 유지
run_dream 성공   explicit count + 완료 시각을 함께 전진
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이때 "성공"은 &lt;code&gt;result is not None&lt;/code&gt;이라는 현재 코드의 비교적 관대한 기준이다. 생성된 memory의 의미적 품질이나 정확성까지 보장하는 quality gate는 아니다.&lt;/p&gt;

&lt;h2&gt;
  
  
  세 번째 경계: 실행 중 상태는 queue가 소유한다
&lt;/h2&gt;

&lt;p&gt;checkpoint를 완료 뒤로 미루면 새 문제가 생긴다. 첫 dream이 아직 실행 중일 때 baseline은 이전 값이므로 scheduler가 같은 collection을 다시 예약할 수 있다.&lt;/p&gt;

&lt;p&gt;Honcho는 이 in-flight window를 metadata timestamp로 흉내 내지 않고 pending &lt;code&gt;QueueItem&lt;/code&gt;으로 확인한다. 같은 collection과 dream type의 작업을 식별하는 work-unit key로 pending dream을 먼저 조회해 두 번째 schedule을 막고, 동시에 insert가 경쟁하더라도 DB partial unique index가 duplicate pending row를 거부한다. 현재 enqueue 경로는 이 충돌을 정상 skip으로 바꾸지 않고 예외를 다시 올리므로, 여기서 확인한 보장은 중복 row가 저장되지 않는다는 범위다. 실패해서 pending 상태가 끝났지만 checkpoint가 전진하지 않았다면 같은 corpus를 다시 시도할 수 있다.&lt;/p&gt;

&lt;p&gt;따라서 세 상태의 owner가 분리된다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;새 입력 eligibility  -&amp;gt; Document.level == explicit
완료된 progress      -&amp;gt; collection dream checkpoint
현재 실행 중         -&amp;gt; pending QueueItem
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h2&gt;
  
  
  테스트가 고정하는 계약
&lt;/h2&gt;

&lt;p&gt;PR #573과 pinned test에는 다음 경계가 명시돼 있다.&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;30 explicit + 40 deductive + 10 inductive는 threshold 50을 통과하지 않는다.&lt;/li&gt;
&lt;li&gt;60 explicit은 threshold를 통과한다.&lt;/li&gt;
&lt;li&gt;100 contradiction + 10 explicit도 통과하지 않는다.&lt;/li&gt;
&lt;li&gt;성공한 dream은 완료 시점의 explicit count와 &lt;code&gt;last_dream_at&lt;/code&gt;을 함께 기록한다.&lt;/li&gt;
&lt;li&gt;pending dream이 있으면 같은 collection의 두 번째 schedule을 막는다.&lt;/li&gt;
&lt;li&gt;
&lt;code&gt;run_dream&lt;/code&gt;이 &lt;code&gt;None&lt;/code&gt;을 반환하면 checkpoint를 전진시키지 않고 같은 corpus의 재시도를 허용한다.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;이번 로컬 환경에서는 Honcho가 요구하는 Python/pytest 조합을 실행하지 못해 테스트 코드를 정적으로 확인했다.&lt;/p&gt;

&lt;h2&gt;
  
  
  다른 pipeline에 옮길 수 있는 계약
&lt;/h2&gt;

&lt;p&gt;이 사례는 AI memory에만 한정되지 않는다. ETL 결과가 다음 ETL의 입력 후보가 되거나, feature가 다시 학습 데이터에 들어가거나, LLM 생성물이 다음 RAG corpus에 포함되는 시스템에서도 같은 질문이 필요하다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;1. source와 derived를 저장 계층뿐 아니라 eligibility에서도 구분한다.
2. progress checkpoint는 enqueue가 아니라 성공 조건 뒤에 전진시킨다.
3. queued/running 상태를 완료 checkpoint에 섞지 않는다.
4. count, session selection, retry가 같은 input cohort를 보게 한다.
5. 파생물 저장 뒤 checkpoint가 실패하면 같은 입력이 재실행될 수 있다. 재실행을 허용한다면
   파생물 write의 deduplication 또는 idempotency를 함께 설계한다.
6. checkpoint 이후에도 남은 side effect가 있다면 그 실패를 다시 처리할 별도 상태를 둔다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;derived data를 영원히 입력에서 제외하라는 뜻은 아니다. 재귀적 개선이 필요한 시스템이라면 어떤 derived level을 언제 다시 입력으로 허용할지 별도의 명시적 정책과 종료 조건을 둬야 한다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Limitations
&lt;/h2&gt;

&lt;ul&gt;
&lt;li&gt;Honcho의 공개 코드와 변경 이력을 분석한 글이며, Honcho를 직접 구현·운영하거나 이 변경에 기여한 경험이 아니다.&lt;/li&gt;
&lt;li&gt;분석 기준은 &lt;code&gt;85239a69b262c944de3c35900b91c88ba9b84f1a&lt;/code&gt;로 고정했다.&lt;/li&gt;
&lt;li&gt;production workload에서 반복 schedule, 비용 증가, 무한 loop를 재현하지 않았다.&lt;/li&gt;
&lt;li&gt;PR #573은 feedback eligibility뿐 아니라 time guard, metadata write, queue coherence를 함께 수정한다.&lt;/li&gt;
&lt;li&gt;non-null &lt;code&gt;DreamResult&lt;/code&gt; 이후 checkpoint하는 현재 기준이 memory 품질까지 검증한다고 주장하지 않는다.&lt;/li&gt;
&lt;li&gt;
&lt;code&gt;last_dream_document_count&lt;/code&gt;가 전체 level count에서 explicit-only count로 의미가 바뀔 때 기존 collection metadata를 어떻게 전환하는지는 이번 분석에서 확인하지 않았다.&lt;/li&gt;
&lt;li&gt;테스트 파일은 확인했지만 현재 로컬 환경에서는 실행하지 못했다.&lt;/li&gt;
&lt;/ul&gt;

&lt;h2&gt;
  
  
  Sources
&lt;/h2&gt;

&lt;ul&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho" rel="noopener noreferrer"&gt;Honcho repository&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/pull/573" rel="noopener noreferrer"&gt;PR #573: threshold and time-guard semantics&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/commit/f37338b855d9fe1ab06e7e4b8e676e6fd01baa47" rel="noopener noreferrer"&gt;PR #573 merge commit&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/a05c2f8ec1f03893d9805a39a30a15f2f061c2f0/src/dreamer/dream_scheduler.py#L229-L335" rel="noopener noreferrer"&gt;Code immediately before the change&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/src/dreamer/dream_scheduler.py#L248-L406" rel="noopener noreferrer"&gt;Pinned explicit-only threshold and pending-queue gate&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/src/dreamer/orchestrator.py#L335-L411" rel="noopener noreferrer"&gt;Pinned success checkpoint&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/src/deriver/enqueue.py#L445-L529" rel="noopener noreferrer"&gt;Pinned enqueue boundary&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/src/models.py#L515-L528" rel="noopener noreferrer"&gt;Pinned pending-dream partial unique index&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/tests/dreamer/test_dream_scheduler.py#L288-L421" rel="noopener noreferrer"&gt;Pinned threshold regression tests&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/tests/dreamer/test_dreamer_integration.py#L175-L245" rel="noopener noreferrer"&gt;Pinned checkpoint and in-flight coherence tests&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/tests/dreamer/test_dreamer_integration.py#L434-L599" rel="noopener noreferrer"&gt;Pinned queue/retry coherence tests&lt;/a&gt;&lt;/li&gt;
&lt;/ul&gt;

</description>
      <category>ai</category>
      <category>dataplatform</category>
      <category>reliability</category>
      <category>opensource</category>
    </item>
    <item>
      <title>느린 LLM 호출 중 DB connection을 잡지 않는 이유</title>
      <dc:creator>박준현</dc:creator>
      <pubDate>Thu, 16 Jul 2026 00:47:00 +0000</pubDate>
      <link>https://dev.to/junhyun-dev/neurin-llm-hocul-jung-db-connectioneul-jabji-anhneun-iyu-3abg</link>
      <guid>https://dev.to/junhyun-dev/neurin-llm-hocul-jung-db-connectioneul-jabji-anhneun-iyu-3abg</guid>
      <description>&lt;h1&gt;
  
  
  느린 LLM 호출 중 DB connection을 잡지 않는 이유
&lt;/h1&gt;

&lt;p&gt;AI 기능의 latency는 모델 응답 시간으로만 끝나지 않는다. 요청이 LLM이나 embedding API를 기다리는 동안 데이터베이스 session까지 긴 범위로 유지하면, 느린 외부 호출이 DB connection pool의 압력으로 전파될 수 있다.&lt;/p&gt;

&lt;p&gt;이 글은 AI memory OSS인 Honcho의 변경 이력과 pinned source를 읽으면서, 외부 호출과 DB transaction/session 경계를 어떻게 분리했는지 추적한 기록이다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Scenario
&lt;/h2&gt;

&lt;p&gt;Honcho의 dialectic 경로(저장된 memory를 근거로 사용자 질문에 답하는 질의 경로)는 답하기 전에 다음 작업을 수행한다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;peer·session·workspace 확인
-&amp;gt; 관련 memory 검색
-&amp;gt; embedding·LLM 호출
-&amp;gt; 필요하면 tool로 추가 조회·기록
-&amp;gt; 답변 생성
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;여기에는 짧은 DB 조회와 상대적으로 느리고 변동성이 큰 외부 호출이 섞여 있다. 두 작업을 하나의 session scope로 묶으면, DB가 필요하지 않은 대기 시간까지 session lifetime에 포함된다.&lt;/p&gt;

&lt;h2&gt;
  
  
  변경 전에는 무엇이 묶여 있었나
&lt;/h2&gt;

&lt;p&gt;PR #477 직전의 &lt;code&gt;agentic_chat&lt;/code&gt;은 하나의 &lt;code&gt;tracked_db&lt;/code&gt; context 안에서 peer와 설정을 읽고, 그 session을 &lt;code&gt;DialecticAgent&lt;/code&gt;에 전달한 뒤, &lt;code&gt;agent.answer()&lt;/code&gt;가 끝날 때까지 같은 context를 유지했다. streaming 경로도 같은 형태였다.&lt;/p&gt;

&lt;p&gt;구조를 단순화하면 다음과 같다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;DB session open
  -&amp;gt; preflight read
  -&amp;gt; DialecticAgent receives session
  -&amp;gt; embedding / memory tools / LLM answer
DB session close
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;commit &lt;code&gt;0533c6d&lt;/code&gt;의 제목도 이 문제를 &lt;code&gt;dialectic held connection&lt;/code&gt;으로 기록한다. 다만 이번 분석에서는 실제 pool checkout 시간이나 장애를 재현하지 않았다. 여기서 확인한 것은 코드의 session scope와 변경 의도다.&lt;/p&gt;

&lt;p&gt;SQLAlchemy의 session 객체를 만들었다고 곧바로 connection을 점유하는 것은 아니다. 하지만 이 경로처럼 SQL을 실행해 transaction이 시작되면 session은 pool에서 빌린 connection을 commit·rollback까지 유지한다. Honcho의 &lt;code&gt;tracked_db&lt;/code&gt;는 종료할 때 &lt;code&gt;rollback()&lt;/code&gt;과 &lt;code&gt;close()&lt;/code&gt;를 호출하므로, SQL 실행 뒤 이 context를 LLM 대기까지 유지하던 범위를 줄이는 것은 이 경로의 connection 점유 구간도 줄이는 일이다.&lt;/p&gt;

&lt;h2&gt;
  
  
  어떻게 경계를 줄였나
&lt;/h2&gt;

&lt;p&gt;변경 후에는 &lt;code&gt;tracked_db("dialectic.preflight")&lt;/code&gt;가 본 작업 전 검증과 설정 조회(preflight), 즉 peer 존재 여부, session/workspace 설정, peer card를 읽는 구간만 감싼다. context가 끝난 다음에 &lt;code&gt;DialecticAgent&lt;/code&gt;를 만들고 LLM 답변을 생성한다. Agent 생성자에서도 DB session 인자가 제거됐다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;short DB preflight
  -&amp;gt; 필요한 값 읽기
DB session close

agent execution
  -&amp;gt; embedding / LLM / tools

tool needs DB
  -&amp;gt; tool-owned short DB session
  -&amp;gt; close
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;핵심은 DB 사용을 없앤 것이 아니다. &lt;strong&gt;요청 전체가 session을 소유하는 대신, DB가 필요한 작업이 자기 범위의 session을 소유하도록 바꾼 것&lt;/strong&gt;이다.&lt;/p&gt;

&lt;h2&gt;
  
  
  pinned current code에서도 유지되는가
&lt;/h2&gt;

&lt;p&gt;분석 기준 commit &lt;code&gt;85239a6&lt;/code&gt;에서도 이 경계는 유지되고 더 구체화돼 있다.&lt;/p&gt;

&lt;ol&gt;
&lt;li&gt;
&lt;code&gt;src/dialectic/chat.py&lt;/code&gt;의 일반·streaming 경로 모두 preflight context를 닫은 뒤 agent를 실행한다.&lt;/li&gt;
&lt;li&gt;
&lt;code&gt;DialecticAgent&lt;/code&gt;는 DB session을 필드로 받지 않는다.&lt;/li&gt;
&lt;li&gt;observation prefetch는 embedding을 먼저 계산한 뒤 &lt;code&gt;search_memory&lt;/code&gt;를 호출한다.&lt;/li&gt;
&lt;li&gt;
&lt;code&gt;search_memory&lt;/code&gt;는 &lt;code&gt;query_documents(db=None, ...)&lt;/code&gt;를 사용해 session lifetime을 하위 함수에 맡긴다.&lt;/li&gt;
&lt;li&gt;external vector store 경로는 외부 검색으로 document ID를 얻은 후에만 짧은 DB session을 열어 row를 가져온다.&lt;/li&gt;
&lt;li&gt;PostgreSQL 안에서 vector 검색을 수행하는 pgvector 경로는 검색 자체가 DB 연산이므로 그 DB query 범위에는 session을 사용한다.&lt;/li&gt;
&lt;li&gt;읽기·쓰기 tool은 필요할 때 각자의 &lt;code&gt;tracked_db&lt;/code&gt; context를 연다.&lt;/li&gt;
&lt;/ol&gt;

&lt;p&gt;따라서 transferable contract는 “외부 호출이 있으면 DB를 절대 쓰지 말라”가 아니다.&lt;/p&gt;

&lt;blockquote&gt;
&lt;p&gt;transaction consistency가 필요한 구간과 외부 네트워크 대기 구간을 분리하고, 각 DB 작업이 필요한 최소 session boundary를 소유하게 한다.&lt;/p&gt;
&lt;/blockquote&gt;

&lt;h2&gt;
  
  
  테스트에서 확인할 수 있는 것
&lt;/h2&gt;

&lt;p&gt;PR #477은 테스트의 tool context에서도 공유 DB session을 제거했다. 테스트 fixture는 tool handler가 여는 독립 session에서 데이터를 볼 수 있도록 준비 데이터를 commit하고, &lt;code&gt;tracked_db&lt;/code&gt;를 fresh session factory로 대체한다. 이는 tool이 요청 전체의 session에 의존하지 않는다는 계약을 테스트 구조에 반영한 것이다.&lt;/p&gt;

&lt;p&gt;이후 pinned code의 message search 테스트에는 external semantic lookup이 &lt;code&gt;tracked_db&lt;/code&gt; 진입보다 먼저 발생하는지 call order로 확인하는 사례도 있다. 반면 이번 로컬 환경에서는 Honcho의 pytest dependency를 실행할 수 없어 테스트 코드를 정적으로만 확인했다.&lt;/p&gt;

&lt;h2&gt;
  
  
  적용할 때 같이 봐야 하는 위험
&lt;/h2&gt;

&lt;p&gt;session scope를 줄이는 것에도 비용이 있다.&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;session 종료 뒤 ORM lazy attribute에 접근하면 detached object 문제가 생길 수 있다.&lt;/li&gt;
&lt;li&gt;하나의 transaction으로 보호해야 하는 read-modify-write를 무리하게 나누면 consistency가 깨질 수 있다.&lt;/li&gt;
&lt;li&gt;preflight와 이후 tool은 서로 다른 transaction snapshot을 볼 수 있다. 경계 사이의 데이터 변경을 허용해도 되는지 별도로 판단해야 한다.&lt;/li&gt;
&lt;li&gt;tool마다 session을 열면 session 수는 짧아지지만 호출 횟수와 transaction 설계는 별도로 검토해야 한다.&lt;/li&gt;
&lt;li&gt;pgvector처럼 외부 호출이 아니라 DB 자체에서 검색하는 경로는 DB session이 필요하다.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;그래서 순서는 “무조건 session을 빨리 닫기”가 아니라 다음과 같다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;1. 실제 consistency boundary를 정한다.
2. 외부 호출 전에 필요한 값을 미리 읽어 변수로 확보한다(materialize).
3. DB scope를 닫는다.
4. 느린 외부 작업을 수행한다.
5. 결과 저장이 필요하면 새롭고 짧은 write scope를 연다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h2&gt;
  
  
  Limitations
&lt;/h2&gt;

&lt;ul&gt;
&lt;li&gt;Honcho의 공개 코드와 변경 이력을 분석한 글이며, Honcho를 직접 구현·운영하거나 이 변경에 기여한 경험이 아니다.&lt;/li&gt;
&lt;li&gt;분석 기준은 &lt;code&gt;85239a69b262c944de3c35900b91c88ba9b84f1a&lt;/code&gt;로 고정했다.&lt;/li&gt;
&lt;li&gt;production 부하, connection pool 고갈, connection checkout 시간을 재현하거나 benchmark하지 않았다.&lt;/li&gt;
&lt;li&gt;테스트 파일은 확인했지만 현재 로컬 환경에는 요구 Python/pytest 조합이 없어 실행하지 못했다.&lt;/li&gt;
&lt;li&gt;감사한 dialectic·memory search 경로 밖의 모든 외부 호출이 같은 계약을 지킨다고 주장하지 않는다.&lt;/li&gt;
&lt;/ul&gt;

&lt;h2&gt;
  
  
  Sources
&lt;/h2&gt;

&lt;ul&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho" rel="noopener noreferrer"&gt;Honcho repository&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/commit/0533c6dd26d2fb4928eae0c837efae608e6351d3" rel="noopener noreferrer"&gt;PR #477 change commit&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/cc3483bfaa92277e349b1f7f1a6b0767c2512085/src/dialectic/chat.py#L42-L83" rel="noopener noreferrer"&gt;Code immediately before the change&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/src/dialectic/chat.py#L42-L78" rel="noopener noreferrer"&gt;Pinned dialectic preflight boundary&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/src/dependencies.py#L38-L68" rel="noopener noreferrer"&gt;Pinned &lt;code&gt;tracked_db&lt;/code&gt; rollback/close boundary&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/src/dialectic/core.py#L151-L205" rel="noopener noreferrer"&gt;Pinned observation prefetch call site&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/src/utils/agent_tools.py#L1062-L1106" rel="noopener noreferrer"&gt;Pinned memory search boundary&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/src/crud/document.py#L316-L421" rel="noopener noreferrer"&gt;Pinned external-vector/DB two-phase query&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/tests/utils/test_agent_tools.py#L118-L132" rel="noopener noreferrer"&gt;Pinned independent-session test setup&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://github.com/plastic-labs/honcho/blob/85239a69b262c944de3c35900b91c88ba9b84f1a/tests/integration/test_message_embeddings.py#L405-L501" rel="noopener noreferrer"&gt;Pinned external-before-DB order test&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://docs.sqlalchemy.org/en/20/orm/session_basics.html" rel="noopener noreferrer"&gt;SQLAlchemy Session basics&lt;/a&gt;&lt;/li&gt;
&lt;/ul&gt;

</description>
      <category>ai</category>
      <category>backend</category>
      <category>python</category>
      <category>opensource</category>
    </item>
    <item>
      <title>skip에서 partition overwrite로: business_date 재처리를 Iceberg로 다시 표현하기</title>
      <dc:creator>박준현</dc:creator>
      <pubDate>Sun, 12 Jul 2026 15:49:01 +0000</pubDate>
      <link>https://dev.to/junhyun-dev/skipeseo-partition-overwritero-businessdate-jaeceorireul-icebergro-dasi-pyohyeonhagi-194i</link>
      <guid>https://dev.to/junhyun-dev/skipeseo-partition-overwritero-businessdate-jaeceorireul-icebergro-dasi-pyohyeonhagi-194i</guid>
      <description>&lt;h1&gt;
  
  
  skip에서 partition overwrite로: business_date 재처리를 Iceberg로 다시 표현하기
&lt;/h1&gt;

&lt;p&gt;이전 글에서는 같은 &lt;code&gt;source_hash&lt;/code&gt;가 다시 들어왔을 때 기존 successful run을 재사용하는 idempotency를 다뤘다.&lt;/p&gt;

&lt;p&gt;하지만 재처리에는 두 종류가 있다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;1. 같은 입력이 다시 들어온 경우
   -&amp;gt; skip이 맞다.

2. 같은 business_date의 정정 입력이 들어온 경우
   -&amp;gt; skip하면 안 된다.
   -&amp;gt; 같은 날짜의 gold 결과를 중복 없이 교체해야 한다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;&lt;code&gt;manufacturing-data-platform-mini&lt;/code&gt;의 B5 slice는 두 번째 문제를 아주 작게 다룬다.&lt;/p&gt;

&lt;p&gt;전체 Spark pipeline을 만든 것이 아니다. &lt;code&gt;gold_daily_metrics&lt;/code&gt; Iceberg table 하나를 local Spark에서 만들고, &lt;code&gt;business_date&lt;/code&gt; partition overwrite와 snapshot evidence만 검증했다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Scenario
&lt;/h2&gt;

&lt;p&gt;이미 아래 gold row가 있다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;business_date=2026-06-29
plant-a / line-1 / gearbox-a
units_produced=120
defect_count=3
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;나중에 같은 &lt;code&gt;business_date=2026-06-29&lt;/code&gt;에 대한 정정 source가 들어온다.&lt;/p&gt;

&lt;p&gt;운영자가 원하는 것은 append가 아니다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;원하지 않는 상태:
  2026-06-29 old row
  2026-06-29 corrected row
  -&amp;gt; 같은 날짜 결과가 중복됨

원하는 상태:
  2026-06-29 corrected row만 남음
  2026-06-30 같은 다른 날짜 partition은 그대로 유지됨
  재처리 전후 snapshot evidence가 남음
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;그래서 이 slice의 질문은 이렇다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;같은 business_date의 정정 source를 처리할 때,
gold table에서 해당 날짜 partition만 중복 없이 교체하고,
어떤 run이 어떤 Iceberg snapshot을 만들었는지 남길 수 있는가?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h2&gt;
  
  
  Decision Pressure
&lt;/h2&gt;

&lt;p&gt;Slice1의 CSV pipeline은 already-successful source를 안전하게 skip할 수 있다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;dataset_id + business_date + source_hash
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이 key가 같으면 같은 입력이다. 다시 계산해도 같은 결과이므로 기존 run을 재사용한다.&lt;/p&gt;

&lt;p&gt;하지만 &lt;code&gt;source_hash&lt;/code&gt;가 달라졌다면 의미가 다르다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;same business_date
different source_hash
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이건 retry가 아니라 correction이다.&lt;/p&gt;

&lt;p&gt;CSV run-folder 방식에서는 새 run output을 만들 수는 있지만, "현재 gold table에서 해당 날짜를 원자적으로 교체한다"는 table-level 의미가 약하다.&lt;/p&gt;

&lt;p&gt;Iceberg를 붙이는 이유는 여기 있다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;source_hash
  -&amp;gt; 같은 입력인지 판단하는 idempotency key

business_date partition
  -&amp;gt; 정정 시 교체할 gold table 범위

snapshot_id
  -&amp;gt; table commit의 evidence
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;즉 Spark/Iceberg는 도구 이름을 추가하려고 붙인 것이 아니라, 재처리 상태 전이를 더 명확히 표현하기 위해 붙였다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Options
&lt;/h2&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Option&lt;/th&gt;
&lt;th&gt;장점&lt;/th&gt;
&lt;th&gt;문제&lt;/th&gt;
&lt;th&gt;판단&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;same source면 항상 재계산&lt;/td&gt;
&lt;td&gt;단순함&lt;/td&gt;
&lt;td&gt;retry 때 불필요한 commit이 계속 생김&lt;/td&gt;
&lt;td&gt;제외&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;corrected source를 append&lt;/td&gt;
&lt;td&gt;구현 쉬움&lt;/td&gt;
&lt;td&gt;같은 날짜 gold row가 중복될 수 있음&lt;/td&gt;
&lt;td&gt;제외&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;whole-table overwrite&lt;/td&gt;
&lt;td&gt;단순함&lt;/td&gt;
&lt;td&gt;다른 날짜 partition까지 지울 위험&lt;/td&gt;
&lt;td&gt;제외&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;
&lt;code&gt;business_date&lt;/code&gt; partition overwrite&lt;/td&gt;
&lt;td&gt;correction 범위가 명확함&lt;/td&gt;
&lt;td&gt;Spark/Iceberg 설정과 test가 필요&lt;/td&gt;
&lt;td&gt;선택&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;MERGE/upsert&lt;/td&gt;
&lt;td&gt;강력함&lt;/td&gt;
&lt;td&gt;이번 skeleton에 과함&lt;/td&gt;
&lt;td&gt;backlog&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;이번 구현은 &lt;code&gt;DataFrameWriterV2.overwritePartitions()&lt;/code&gt;를 사용했다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight python"&gt;&lt;code&gt;&lt;span class="n"&gt;corrected_df&lt;/span&gt;&lt;span class="p"&gt;.&lt;/span&gt;&lt;span class="nf"&gt;writeTo&lt;/span&gt;&lt;span class="p"&gt;(&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="s"&gt;local.db.gold_daily_metrics&lt;/span&gt;&lt;span class="sh"&gt;"&lt;/span&gt;&lt;span class="p"&gt;).&lt;/span&gt;&lt;span class="nf"&gt;overwritePartitions&lt;/span&gt;&lt;span class="p"&gt;()&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;SQL &lt;code&gt;INSERT OVERWRITE&lt;/code&gt;를 바로 쓰지 않은 이유는, 설정을 잘못 잡으면 static overwrite처럼 동작해 전체 table을 덮는 실수를 놓칠 수 있기 때문이다.&lt;/p&gt;

&lt;p&gt;따라서 핵심 test는 단순히 "정정 날짜가 바뀌었나"가 아니라 이것이다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;D partition은 corrected rows로 교체된다.
D2 partition은 그대로 남는다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h2&gt;
  
  
  Decision
&lt;/h2&gt;

&lt;p&gt;이번 slice는 local Spark/Iceberg walking skeleton으로 고정했다.&lt;/p&gt;

&lt;p&gt;버전 pin:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;Python: 3.10.12
Java: OpenJDK 17.0.19
PySpark: 3.5.8
Iceberg: 1.11.0
runtime jar: org.apache.iceberg:iceberg-spark-runtime-3.5_2.12:1.11.0
catalog: local hadoop catalog
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;구현 범위:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;1. local SparkSession 생성
2. Iceberg hadoop catalog 설정
3. local.db.gold_daily_metrics table 생성
4. initial rows append
5. same source_hash retry는 skip처럼 처리해서 새 snapshot 없음
6. corrected rows로 business_date partition overwrite
7. current table rows 확인
8. snapshots metadata 확인
9. run_id -&amp;gt; snapshot_id evidence JSON 저장
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;테이블은 하나다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;local.db.gold_daily_metrics
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;partition column도 하나다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;business_date
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h2&gt;
  
  
  Evidence
&lt;/h2&gt;

&lt;p&gt;구현 evidence:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;requirements-spark.txt
src/manufacturing_data_platform/pipeline/spark_iceberg_skeleton.py
tests/test_spark_iceberg_skeleton.py
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;검증 로그:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;2026-07-11 — Spark/Iceberg single-gold-table walking skeleton
pytest tests/test_spark_iceberg_skeleton.py -q: 2 passed
pytest: 40 passed
Spark/Iceberg CLI: passed
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;실행 명령:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight shell"&gt;&lt;code&gt;pip &lt;span class="nb"&gt;install&lt;/span&gt; &lt;span class="nt"&gt;-r&lt;/span&gt; requirements-spark.txt

&lt;span class="nv"&gt;PYTHONPATH&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;src python &lt;span class="nt"&gt;-m&lt;/span&gt; manufacturing_data_platform.pipeline.spark_iceberg_skeleton &lt;span class="se"&gt;\&lt;/span&gt;
  &lt;span class="nt"&gt;--warehouse&lt;/span&gt; /tmp/manufacturing-mini-iceberg-warehouse &lt;span class="se"&gt;\&lt;/span&gt;
  &lt;span class="nt"&gt;--output-dir&lt;/span&gt; /tmp/manufacturing-mini-iceberg-evidence &lt;span class="se"&gt;\&lt;/span&gt;
  &lt;span class="nt"&gt;--clean&lt;/span&gt;
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;실제 evidence 일부:&lt;/p&gt;

&lt;p&gt;&lt;code&gt;snapshot_id&lt;/code&gt; 값은 Iceberg commit 때 생성되는 값이라 실행마다 달라질 수 있다. 중요한 것은 id 숫자 자체가 아니라, same-source retry에서는 새 snapshot이 없고 correction에서는 새 snapshot이 생긴다는 상태 변화다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight json"&gt;&lt;code&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"dataset_id"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"manufacturing_daily_metrics"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"table"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"local.db.gold_daily_metrics"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"business_date"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"2026-06-29"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"runs"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"run_id"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"spark-skeleton-r1"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"source_hash"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"source-hash-initial-001"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"status"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"processed"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"gold_snapshot_id"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;2920694863405545739&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"snapshot_count"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;1&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;},&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"run_id"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"spark-skeleton-r1-retry"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"source_hash"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"source-hash-initial-001"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"status"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"skipped"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"gold_snapshot_id"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;2920694863405545739&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"snapshot_count"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;1&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;},&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"run_id"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"spark-skeleton-r2"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"source_hash"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"source-hash-corrected-002"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"status"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"processed"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"gold_snapshot_id"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;8586499016384598474&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"snapshot_count"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;2&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;}&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="p"&gt;],&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"partition_overwrite_assertions"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"target_partition_row_count"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;1&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"corrected_row_count"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;1&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"snapshot_increment"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;1&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"same_source_created_snapshot"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="kc"&gt;false&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="p"&gt;}&lt;/span&gt;&lt;span class="w"&gt;
&lt;/span&gt;&lt;span class="p"&gt;}&lt;/span&gt;&lt;span class="w"&gt;
&lt;/span&gt;&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;current gold rows도 확인했다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight json"&gt;&lt;code&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"rows"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"business_date"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"2026-06-29"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"plant_id"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"plant-a"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"line_id"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"line-1"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"product_code"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"gearbox-a"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"units_produced"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;150&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"defect_count"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;6&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"defect_rate"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mf"&gt;0.04&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;},&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"business_date"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"2026-06-30"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"plant_id"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"plant-a"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"line_id"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"line-1"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"product_code"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"gearbox-a"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"units_produced"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;50&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"defect_count"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;1&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"defect_rate"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mf"&gt;0.02&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;}&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="p"&gt;]&lt;/span&gt;&lt;span class="w"&gt;
&lt;/span&gt;&lt;span class="p"&gt;}&lt;/span&gt;&lt;span class="w"&gt;
&lt;/span&gt;&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이 결과가 의미하는 것은 단순하다.&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;code&gt;2026-06-29&lt;/code&gt;는 corrected row 하나만 남았다.&lt;/li&gt;
&lt;li&gt;
&lt;code&gt;2026-06-30&lt;/code&gt; partition은 사라지지 않았다.&lt;/li&gt;
&lt;li&gt;같은 &lt;code&gt;source_hash&lt;/code&gt; retry는 새 snapshot을 만들지 않았다.&lt;/li&gt;
&lt;li&gt;corrected source는 새 snapshot을 만들었다.&lt;/li&gt;
&lt;li&gt;
&lt;code&gt;run_id&lt;/code&gt;와 &lt;code&gt;snapshot_id&lt;/code&gt;를 구분해서 evidence로 남겼다.&lt;/li&gt;
&lt;/ul&gt;

&lt;h2&gt;
  
  
  Why Snapshot ID Is Not Run ID
&lt;/h2&gt;

&lt;p&gt;&lt;code&gt;run_id&lt;/code&gt;와 &lt;code&gt;snapshot_id&lt;/code&gt;는 다른 세계의 id다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;run_id:
  pipeline execution identity
  예: spark-skeleton-r2

snapshot_id:
  Iceberg table commit identity
  예: 8586499016384598474
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이번 skeleton에서는 한 run이 gold table commit 하나를 만든다는 invariant를 뒀다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;one pipeline run -&amp;gt; one gold table commit -&amp;gt; one gold_snapshot_id
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;그래서 &lt;code&gt;run_id -&amp;gt; snapshot_id&lt;/code&gt; mapping이 가능하다.&lt;/p&gt;

&lt;p&gt;나중에 한 run이 여러 Iceberg table에 commit하거나, 한 table에 여러 번 commit하면 이 관계는 1:N으로 바뀐다. 그건 backlog다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Limitations
&lt;/h2&gt;

&lt;p&gt;이건 production lakehouse가 아니다.&lt;/p&gt;

&lt;p&gt;명확한 한계:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;single gold table walking skeleton
full bronze/silver/gold Spark rewrite 아님
Spark-based quality suite 아님
MERGE/upsert 아님
Iceberg rollback system 아님
time-travel read demo는 아직 본문 claim으로 쓰지 않음
retention/expire snapshots 아님
concurrent writer handling 아님
Airflow-triggered Spark runtime 아님
cluster/Kubernetes runtime 아님
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;그리고 이 구현은 local proof다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;local SparkSession
local hadoop catalog
/tmp warehouse
synthetic data
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이 범위를 넘겨서 "운영 lakehouse를 구축했다"고 말하면 과장이다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Reference Links
&lt;/h2&gt;

&lt;p&gt;이번 구현에서 참고한 공식 기준:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;&lt;a href="https://spark.apache.org/docs/3.5.8/" rel="noopener noreferrer"&gt;Apache Spark 3.5.8 docs&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://iceberg.apache.org/docs/latest/spark-getting-started/" rel="noopener noreferrer"&gt;Apache Iceberg Spark Getting Started&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://iceberg.apache.org/multi-engine-support/" rel="noopener noreferrer"&gt;Apache Iceberg Multi-Engine Support&lt;/a&gt;&lt;/li&gt;
&lt;li&gt;&lt;a href="https://central.sonatype.com/artifact/org.apache.iceberg/iceberg-spark-runtime-3.5_2.12" rel="noopener noreferrer"&gt;Maven Central: iceberg-spark-runtime-3.5_2.12&lt;/a&gt;&lt;/li&gt;
&lt;/ul&gt;

&lt;h2&gt;
  
  
  정리
&lt;/h2&gt;

&lt;p&gt;same-source retry와 corrected-source rerun은 다르게 다뤄야 한다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;same source_hash
  -&amp;gt; skip
  -&amp;gt; no new snapshot

different source_hash + same business_date
  -&amp;gt; correction
  -&amp;gt; business_date partition overwrite
  -&amp;gt; new snapshot evidence
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이 작은 skeleton은 Spark/Iceberg 전체 이식이 아니라, 이 상태 전이가 실제로 가능한지 검증한 것이다.&lt;/p&gt;

&lt;p&gt;코드: &lt;a href="https://github.com/junhyun-dev/manufacturing-data-platform-mini" rel="noopener noreferrer"&gt;github.com/junhyun-dev/manufacturing-data-platform-mini&lt;/a&gt;&lt;/p&gt;

&lt;h2&gt;
  
  
  정리
&lt;/h2&gt;

&lt;p&gt;&lt;code&gt;source_hash&lt;/code&gt; skip만으로는 정정 파일을 설명하기 어렵다. 같은 입력은 skip해야 하지만, 같은 &lt;code&gt;business_date&lt;/code&gt;의 다른 입력은 기존 결과를 중복 없이 교체해야 한다.&lt;/p&gt;

&lt;p&gt;이 slice에서는 단일 local Iceberg gold table에서 그 상태 전이를 검증했다. 같은 &lt;code&gt;source_hash&lt;/code&gt; rerun은 새 snapshot을 만들지 않고, 다른 &lt;code&gt;source_hash&lt;/code&gt; correction은 &lt;code&gt;DataFrameWriterV2.overwritePartitions()&lt;/code&gt;로 해당 &lt;code&gt;business_date&lt;/code&gt; partition만 교체한다. 다른 날짜 partition 보존, 중복 row 부재, &lt;code&gt;run_id&lt;/code&gt;와 &lt;code&gt;snapshot_id&lt;/code&gt; 분리를 테스트로 확인했다.&lt;/p&gt;

&lt;p&gt;범위는 명확하다. 이것은 production lakehouse나 full Spark medallion rewrite가 아니라, correction rerun을 Iceberg partition overwrite로 표현할 수 있는지 확인한 local walking skeleton이다.&lt;/p&gt;

</description>
      <category>dataengineering</category>
      <category>python</category>
      <category>etl</category>
      <category>learning</category>
    </item>
    <item>
      <title>wide CSV 여러 개를 EAV로 모아 gold mart 만들기</title>
      <dc:creator>박준현</dc:creator>
      <pubDate>Sun, 12 Jul 2026 15:49:00 +0000</pubDate>
      <link>https://dev.to/junhyun-dev/wide-csv-yeoreo-gaereul-eavro-moa-gold-mart-mandeulgi-4799</link>
      <guid>https://dev.to/junhyun-dev/wide-csv-yeoreo-gaereul-eavro-moa-gold-mart-mandeulgi-4799</guid>
      <description>&lt;h1&gt;
  
  
  wide CSV 여러 개를 EAV로 모아 gold mart 만들기
&lt;/h1&gt;

&lt;p&gt;현실의 데이터 소스는 한 가지 모양으로 오지 않는다.&lt;/p&gt;

&lt;p&gt;같은 의미의 값도 어떤 파일에서는 &lt;code&gt;생산수량&lt;/code&gt;, 다른 파일에서는 &lt;code&gt;units&lt;/code&gt;, 또 다른 파일에서는 &lt;code&gt;made&lt;/code&gt;로 올 수 있다. 온도도 어떤 곳은 섭씨, 어떤 곳은 화씨일 수 있다. 이걸 매번 pipeline code에 &lt;code&gt;if source == ...&lt;/code&gt;로 박기 시작하면 source가 늘 때마다 코드가 지저분해진다.&lt;/p&gt;

&lt;p&gt;&lt;code&gt;manufacturing-data-platform-mini&lt;/code&gt;의 EAV mini slice는 이 문제를 작게 다룬다. 여러 wide CSV를 mapping config로 표준 attribute에 맞춘 뒤, EAV long format으로 모으고, 다시 gold metric mart로 pivot/aggregate한다. 데이터는 모두 synthetic이고, 회사 코드·고객 데이터·실제 schema는 쓰지 않았다.&lt;/p&gt;

&lt;h2&gt;
  
  
  1. Scenario
&lt;/h2&gt;

&lt;p&gt;서로 다른 공장/라인/벤더에서 비슷한 제조 지표 파일이 들어온다.&lt;/p&gt;

&lt;p&gt;예:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;plant_a.csv:
  설비ID, 생산수량, 불량수, 온도C, 압력kPa

plant_b.csv:
  machine_id, output_units, defects, temp_f, pressure_bar

vendor_d.csv:
  unit_name, made, scrap, deg_c, kpa
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;비즈니스적으로는 같은 지표를 보고 싶다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;units_produced
defect_count
temperature_c
pressure_kpa
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;문제는 source마다 column name과 unit이 다르다는 점이다.&lt;/p&gt;

&lt;h2&gt;
  
  
  2. Decision Pressure
&lt;/h2&gt;

&lt;p&gt;단순 구현은 source마다 코드를 늘린다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;if source == "plant_a":
  생산수량을 units_produced로 읽는다

if source == "plant_b":
  output_units를 units_produced로 읽는다
  temp_f를 섭씨로 변환한다

if source == "vendor_d":
  made를 units_produced로 읽는다
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이 방식은 작게는 빨라 보이지만 source가 늘수록 문제가 된다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;새 파일 형식마다 pipeline code를 고쳐야 한다.
column mapping과 transform logic이 섞인다.
unit conversion이 흩어진다.
quality check가 source별로 갈라진다.
gold mart grain을 설명하기 어려워진다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;그래서 mapping은 config로 빼고, pipeline은 표준 attribute를 처리하게 만들었다.&lt;/p&gt;

&lt;h2&gt;
  
  
  3. Options
&lt;/h2&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;option&lt;/th&gt;
&lt;th&gt;result&lt;/th&gt;
&lt;th&gt;risk&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;source별 hard-coded parser&lt;/td&gt;
&lt;td&gt;처음엔 빠름&lt;/td&gt;
&lt;td&gt;source가 늘 때 code change 반복&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;모든 source를 wide table 하나로 합치기&lt;/td&gt;
&lt;td&gt;보기 쉬움&lt;/td&gt;
&lt;td&gt;sparse/heterogeneous column 폭발&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;EAV long format&lt;/td&gt;
&lt;td&gt;이종 attribute를 표준 형태로 모음&lt;/td&gt;
&lt;td&gt;pivot/quality 설계가 필요&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;full mapping DSL/rules engine&lt;/td&gt;
&lt;td&gt;유연함&lt;/td&gt;
&lt;td&gt;mini project에는 과함&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;이 프로젝트의 선택은 단순한 JSON mapping + EAV long + gold pivot이다.&lt;/p&gt;

&lt;h2&gt;
  
  
  4. Decision
&lt;/h2&gt;

&lt;p&gt;각 source는 JSON config로 자신의 column을 표준 attribute에 매핑한다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;source column -&amp;gt; standard attribute
output_units  -&amp;gt; units_produced
temp_f        -&amp;gt; temperature_c with f_to_c
pressure_bar  -&amp;gt; pressure_kpa with bar_to_kpa
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;pipeline 흐름:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;wide CSVs
-&amp;gt; mapping configs
-&amp;gt; EAV long rows
-&amp;gt; gold entity_daily_metrics
-&amp;gt; quality checks
-&amp;gt; catalog/lineage
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;EAV row의 핵심 shape:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;entity_id
business_date
attribute
value
value_type
source_id
source_file_id
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;&lt;code&gt;source_file_id&lt;/code&gt;는 file content hash다. 즉 EAV slice도 file-level idempotency를 갖는다.&lt;/p&gt;

&lt;p&gt;더 정확히는 &lt;code&gt;source_file_id&lt;/code&gt;가 각 source file의 hash이고, run-level idempotency key는 모든 source file hash를 합친 &lt;code&gt;source_hash&lt;/code&gt;다. 그래서 같은 source 묶음을 같은 &lt;code&gt;business_date&lt;/code&gt;로 다시 처리하면 기존 successful run을 재사용한다.&lt;/p&gt;

&lt;h2&gt;
  
  
  5. Evidence
&lt;/h2&gt;

&lt;p&gt;관련 코드와 config:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;src/manufacturing_data_platform/pipeline/eav.py
src/manufacturing_data_platform/pipeline/run_eav.py
src/manufacturing_data_platform/pipeline/sample_eav.py
config/eav_mappings/plant_a.json
config/eav_mappings/plant_b.json
config/eav_mappings/line_c.json
tests/test_eav_pipeline.py
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;핵심 테스트:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;test_transform_to_eav_maps_and_converts_units
test_transform_to_eav_captures_type_errors_gracefully
test_transform_eav_to_gold_aggregates_sum_and_avg
test_eav_run_passes_and_unifies_three_formats
test_new_format_is_onboarded_by_adding_one_config
test_eav_idempotent_rerun_is_skipped
test_mapping_coverage_fails_when_required_attribute_unmapped
test_unmapped_columns_warn_does_not_fail
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;특히 &lt;code&gt;test_new_format_is_onboarded_by_adding_one_config&lt;/code&gt;는 중요하다.&lt;/p&gt;

&lt;p&gt;이 테스트는 기존 3개 sample source가 있는 상태에서 &lt;code&gt;vendor_d.csv&lt;/code&gt;와 &lt;code&gt;vendor_d.json&lt;/code&gt;을 추가한다. pipeline code는 바꾸지 않는다. 그 다음 run 결과에 새 entity &lt;code&gt;VD-1&lt;/code&gt;이 들어오는지 확인한다.&lt;/p&gt;

&lt;p&gt;즉 claim은 이렇다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;새 파일 형식 하나는 mapping config 추가만으로 온보딩된다.
pipeline code change는 없다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;검증 로그:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;2026-06-30 EAV mini slice:
  pytest: 33 passed
  EAV CLI run 1: status=processed, quality_passed=true
  EAV CLI run 2: status=skipped
  conservation: EAV units 540 == gold units 540

2026-07-08 publication readiness check:
  pytest: 33 passed
  EAV JSON CLI: passed, status=processed, quality_passed=true

2026-07-10 B3 publication evidence check:
  pytest: 35 passed
  EAV JSON CLI run 1: status=processed, quality_passed=true
  EAV JSON CLI run 2: status=skipped
  gold_rows=4, gold_units_total=540, gold_defects_total=12
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;실제 CLI 검증 요약:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight json"&gt;&lt;code&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"run1_status"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"processed"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"run2_status"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"skipped"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"dataset_id"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"manufacturing_wide_eav"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"business_date"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"2026-06-29"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"gold_rows"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;4&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"entities"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="s2"&gt;"EQP-A1"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"EQP-A2"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"LN-C1"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"MC-B1"&lt;/span&gt;&lt;span class="p"&gt;],&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"gold_units_total"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;540&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"gold_defects_total"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;12&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"lineage_layers"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="s2"&gt;"bronze"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"silver_eav"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"gold"&lt;/span&gt;&lt;span class="p"&gt;]&lt;/span&gt;&lt;span class="w"&gt;
&lt;/span&gt;&lt;span class="p"&gt;}&lt;/span&gt;&lt;span class="w"&gt;
&lt;/span&gt;&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h2&gt;
  
  
  6. Quality Checks
&lt;/h2&gt;

&lt;p&gt;EAV slice도 단순 transform으로 끝나지 않는다.&lt;/p&gt;

&lt;p&gt;quality checks:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;mapping_coverage
unmapped_source_columns (warn)
not_null_value
accepted_values_attribute
value_type_valid
numeric_range_within_bounds
eav_to_gold_conservation
freshness_business_date
schema_drift
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;EAV에서 gold로 pivot/aggregate할 때 additive measure가 보존되는지도 확인한다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;EAV units_produced total == gold units_produced total
EAV defect_count total == gold defect_count total
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;그리고 manufacturing slice와 달리, EAV slice는 unparseable value를 바로 crash하지 않는다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;unparseable value
-&amp;gt; value = None
-&amp;gt; value_type_valid = fail
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;즉 bad value를 quality failure로 드러낸다.&lt;/p&gt;

&lt;h2&gt;
  
  
  7. Claim Boundary
&lt;/h2&gt;

&lt;p&gt;정직한 한계:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;데이터는 fully synthetic이다.
회사 코드, 고객 데이터, 실제 schema를 쓰지 않았다.
full mapping DSL이나 UI는 없다.
EAV가 모든 데이터 모델링 문제의 정답이라는 주장이 아니다.
Spark/Iceberg/Kafka는 이 slice에 구현되어 있지 않다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;실무 경험 claim도 조심해서 말해야 한다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;말해도 되는 것:
  실무에서는 EAV 기반 구조를 운영·개선하며 다양한 파일 양식을 처리했다.
  개인 프로젝트에서는 synthetic data로 wide -&amp;gt; EAV -&amp;gt; gold flow를 직접 구현했다.

말하면 안 되는 것:
  실무 EAV 시스템을 내가 처음부터 설계/구현했다.
  회사 schema를 이 프로젝트에 가져왔다.
  real customer data로 검증했다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h2&gt;
  
  
  8. Why This Matters For The Portfolio
&lt;/h2&gt;

&lt;p&gt;이 slice는 단순히 EAV라는 단어를 넣기 위한 것이 아니다.&lt;/p&gt;

&lt;p&gt;보여주는 능력:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;이종 source를 표준 attribute로 harmonize
mapping config와 transform logic 분리
wide -&amp;gt; long -&amp;gt; mart 모델링
unit conversion
data quality check 설계
새 format 온보딩 contract
clean-room boundary 관리
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;즉 데이터 엔지니어링 포트폴리오에서 "파이프라인을 만들었다"보다 한 단계 더 나아가, source 다양성과 모델링 압력을 다뤘다는 증거가 된다.&lt;/p&gt;

&lt;h2&gt;
  
  
  9. 정리
&lt;/h2&gt;

&lt;p&gt;서로 다른 wide CSV를 하나의 mart로 합치려면 컬럼 이름만 맞추는 것으로는 부족하다. 어떤 source column이 어떤 표준 attribute가 되는지, 단위 변환은 어디서 하는지, 새 format은 어떤 계약으로 들어오는지까지 정해야 한다.&lt;/p&gt;

&lt;p&gt;이 slice에서는 synthetic data만 사용해 wide -&amp;gt; EAV -&amp;gt; gold 흐름을 만들었다. 새 파일 형식은 pipeline code 변경 없이 mapping config 하나로 온보딩되고, mapping coverage, value type, conservation check로 결과를 검증한다. 실무 경험과 개인 프로젝트 구현 claim은 분리해서 말한다.&lt;/p&gt;

</description>
      <category>dataengineering</category>
      <category>python</category>
      <category>etl</category>
      <category>learning</category>
    </item>
    <item>
      <title>schema drift를 fail이 아니라 warn으로 둔 이유</title>
      <dc:creator>박준현</dc:creator>
      <pubDate>Sun, 12 Jul 2026 15:48:59 +0000</pubDate>
      <link>https://dev.to/junhyun-dev/schema-driftreul-faili-anira-warneuro-dun-iyu-540h</link>
      <guid>https://dev.to/junhyun-dev/schema-driftreul-faili-anira-warneuro-dun-iyu-540h</guid>
      <description>&lt;h1&gt;
  
  
  schema drift를 fail이 아니라 warn으로 둔 이유
&lt;/h1&gt;

&lt;p&gt;데이터 파이프라인에서 source schema가 바뀌는 순간은 애매하다.&lt;/p&gt;

&lt;p&gt;무조건 무시하면 운영자는 입력 구조가 바뀐 사실을 모른다. 반대로 모든 schema 변화를 실패로 처리하면, 정상적인 컬럼 추가까지 daily run을 막아버린다.&lt;/p&gt;

&lt;p&gt;&lt;code&gt;manufacturing-data-platform-mini&lt;/code&gt;에서는 이 문제를 작게 다뤘다. synthetic manufacturing CSV의 실제 header를 기준으로 &lt;code&gt;schema_hash&lt;/code&gt;를 만들고, previous successful run과 비교해 달라졌으면 &lt;code&gt;schema_drift&lt;/code&gt; quality check를 &lt;code&gt;warn&lt;/code&gt;으로 남긴다. 단, required column이 빠져 silver/gold contract를 만들 수 없는 경우는 현재 &lt;code&gt;ValueError&lt;/code&gt;로 빠르게 실패한다.&lt;/p&gt;

&lt;h2&gt;
  
  
  1. Scenario
&lt;/h2&gt;

&lt;p&gt;어느 날 source CSV에 새 컬럼이 추가된다.&lt;/p&gt;

&lt;p&gt;기존 header:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;event_time,plant_id,line_id,work_order_id,machine_id,product_code,
operation,units_produced,defect_count,cycle_time_ms,business_date
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;새 header:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;event_time,plant_id,line_id,work_order_id,machine_id,product_code,
operation,units_produced,defect_count,cycle_time_ms,business_date,operator_id
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;&lt;code&gt;operator_id&lt;/code&gt;는 아직 silver/gold mart에서 쓰지 않는다. 하지만 source 구조가 바뀐 사실은 기록되어야 한다.&lt;/p&gt;

&lt;h2&gt;
  
  
  2. Decision Pressure
&lt;/h2&gt;

&lt;p&gt;schema drift에서 중요한 질문은 단순히 "바뀌었나?"가 아니다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;바뀐 것을 운영자가 알 수 있는가?
정상적인 컬럼 추가 때문에 pipeline을 멈춰야 하는가?
downstream gold mart contract가 조용히 바뀌지는 않는가?
이전 successful run과 지금 run의 schema identity를 비교할 수 있는가?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;초기 구현에서는 한 가지 실제 버그가 있었다. &lt;code&gt;schema_hash&lt;/code&gt;가 고정된 required column 목록에 너무 묶여 있어서, 추가 컬럼이 들어와도 hash가 바뀌지 않았다. 즉 &lt;code&gt;operator_id&lt;/code&gt;가 추가되어도 drift가 보이지 않았다.&lt;/p&gt;

&lt;p&gt;이 문제를 고치기 위해 &lt;code&gt;read_rows&lt;/code&gt;가 실제 CSV header를 반환하고, 그 실제 header 기준으로 &lt;code&gt;schema_hash&lt;/code&gt;를 계산하도록 바꿨다.&lt;/p&gt;

&lt;h2&gt;
  
  
  3. Options
&lt;/h2&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;option&lt;/th&gt;
&lt;th&gt;result&lt;/th&gt;
&lt;th&gt;risk&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;ignore drift&lt;/td&gt;
&lt;td&gt;pipeline은 계속 돈다&lt;/td&gt;
&lt;td&gt;source 변화가 보이지 않음&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;fail every drift&lt;/td&gt;
&lt;td&gt;변화에 강하게 반응&lt;/td&gt;
&lt;td&gt;정상적인 컬럼 추가도 막음&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;warn and continue&lt;/td&gt;
&lt;td&gt;변화가 보이고 run도 계속됨&lt;/td&gt;
&lt;td&gt;warning을 inspect해야 함&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;auto-evolve silver/gold&lt;/td&gt;
&lt;td&gt;새 컬럼을 바로 사용 가능&lt;/td&gt;
&lt;td&gt;downstream contract가 조용히 바뀔 수 있음&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;full schema registry&lt;/td&gt;
&lt;td&gt;production에 가까움&lt;/td&gt;
&lt;td&gt;mini slice에는 무거움&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;이 프로젝트의 선택은 &lt;code&gt;warn and continue&lt;/code&gt;다.&lt;/p&gt;

&lt;h2&gt;
  
  
  4. Decision
&lt;/h2&gt;

&lt;p&gt;현재 contract는 이렇다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;previous successful run이 없으면:
  schema_drift = pass
  baseline schema established

current schema_hash == previous successful schema_hash:
  schema_drift = pass

current schema_hash != previous successful schema_hash:
  schema_drift = warn
  quality_passed는 true 유지
  run/lineage record에 previous/current schema_hash 저장

required column missing:
  ValueError fast fail
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;중요한 구분:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;schema_drift warn:
  source shape이 바뀌었다. 하지만 현재 silver/gold contract는 만들 수 있다.

missing required column failure:
  현재 mart contract를 만들 수 없다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이 선택은 "schema registry를 만들었다"는 뜻이 아니다. schema change를 invisible하게 두지 않고, quality/check result와 run metadata에 남긴다는 작은 contract다.&lt;/p&gt;

&lt;p&gt;비슷한 방향의 외부 기준도 있다.&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;a href="https://docs.getdbt.com/docs/build/data-tests" rel="noopener noreferrer"&gt;dbt data tests&lt;/a&gt;는 &lt;code&gt;unique&lt;/code&gt;, &lt;code&gt;not_null&lt;/code&gt;, &lt;code&gt;accepted_values&lt;/code&gt; 같은 generic test로 data contract를 드러낸다.&lt;/li&gt;
&lt;li&gt;
&lt;a href="https://docs.getdbt.com/reference/resource-configs/severity" rel="noopener noreferrer"&gt;dbt severity config&lt;/a&gt;는 test 결과를 error 또는 warn으로 다룰 수 있게 한다.&lt;/li&gt;
&lt;li&gt;
&lt;a href="https://docs.greatexpectations.io/docs/cloud/expectations/expectations_overview/" rel="noopener noreferrer"&gt;Great Expectations&lt;/a&gt;는 Expectation을 데이터에 대한 검증 가능한 assertion으로 본다.&lt;/li&gt;
&lt;li&gt;
&lt;a href="https://iceberg.apache.org/docs/latest/evolution/" rel="noopener noreferrer"&gt;Apache Iceberg schema evolution&lt;/a&gt;은 schema 변화를 table metadata 차원에서 다루는 더 큰 방향이다.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;이 프로젝트는 그 도구들을 구현한 것이 아니라, 같은 운영 질문을 mini pipeline의 &lt;code&gt;schema_hash&lt;/code&gt;와 &lt;code&gt;schema_drift&lt;/code&gt; warning으로 작게 연습한다.&lt;/p&gt;

&lt;h2&gt;
  
  
  5. Evidence
&lt;/h2&gt;

&lt;p&gt;관련 코드와 문서:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;src/manufacturing_data_platform/pipeline/lakehouse.py
learn/reference-decisions/schema-drift.md
VERIFICATION_LOG.md
README.md
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;핵심 테스트:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;tests/test_lakehouse_pipeline.py
  test_schema_drift_helper_states
  test_schema_drift_warns_against_previous_successful_run
  test_schema_stable_when_schema_unchanged_across_dates
  test_schema_drift_warns_on_added_column
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;특히 &lt;code&gt;test_schema_drift_warns_on_added_column&lt;/code&gt;은 &lt;code&gt;operator_id&lt;/code&gt; 컬럼이 추가됐을 때:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;schema_drift.status == "warn"
previous schema hash != current schema hash
quality_passed is True
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;를 확인한다.&lt;/p&gt;

&lt;p&gt;검증 로그에도 이 버그 수정이 남아 있다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;2026-06-30 pre-Codex self-audit:
  schema_hash was computed from fixed REQUIRED_COLUMNS, not the actual CSV header.
  Fix: read_rows now returns the actual header; schema_hash uses that actual header.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;그리고 최신 local verification에서도 전체 테스트와 CLI가 통과했다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;2026-07-10:
  pytest: 35 passed
  lakehouse JSON CLI: passed
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h2&gt;
  
  
  6. Limitations
&lt;/h2&gt;

&lt;p&gt;정직한 한계:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;full schema registry는 없다.
column-level diff UI도 없다.
Iceberg/Delta schema evolution은 아직 구현되지 않았다.
schema_drift는 hash 중심이라 어떤 컬럼이 바뀌었는지 사람이 바로 보기엔 제한이 있다.
warning을 운영자가 inspect하지 않으면 놓칠 수 있다.
missing required column은 아직 structured quality report가 아니라 ValueError fast fail이다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;그래서 이 글의 claim은 작게 둔다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;actual CSV header 기준 schema_hash를 만들고,
previous successful run과 비교해 schema drift를 warn으로 드러내는 mini contract를 구현했다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h2&gt;
  
  
  7. Why This Connects To Iceberg Later
&lt;/h2&gt;

&lt;p&gt;Slice1은 schema drift를 detect하고 warn으로 남긴다.&lt;/p&gt;

&lt;p&gt;Iceberg로 가면 질문이 바뀐다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;새 컬럼을 Iceberg table에 add column으로 반영할 것인가?
어떤 변화는 허용하고 어떤 변화는 금지할 것인가?
과거 snapshot은 어떤 schema로 읽히는가?
downstream gold mart contract가 조용히 바뀌지 않게 어떻게 막을 것인가?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;즉 Slice1의 &lt;code&gt;schema_hash + warn&lt;/code&gt;은 나중에 Iceberg schema evolution으로 이어지는 출발점이다. 하지만 현재 repo에서 Iceberg schema evolution은 design-only이며 구현 evidence가 아니다.&lt;/p&gt;

&lt;h2&gt;
  
  
  8. 정리
&lt;/h2&gt;

&lt;p&gt;schema drift를 전부 failure로 처리하면 정상적인 source 확장을 막을 수 있다. 반대로 아무 기록 없이 통과시키면 downstream contract가 조용히 흔들릴 수 있다.&lt;/p&gt;

&lt;p&gt;이 slice에서는 실제 CSV header 기준 &lt;code&gt;schema_hash&lt;/code&gt;를 남기고, previous successful run과 달라진 header를 &lt;code&gt;warn&lt;/code&gt;으로 드러낸다. run은 계속되지만 운영자는 source 구조 변화를 inspect할 수 있다. required column이 빠져 silver/gold contract를 만들 수 없는 경우는 빠르게 실패시킨다.&lt;/p&gt;

</description>
      <category>dataengineering</category>
      <category>python</category>
      <category>etl</category>
      <category>learning</category>
    </item>
    <item>
      <title>source_hash로 같은 입력 재처리를 안전하게 skip하기</title>
      <dc:creator>박준현</dc:creator>
      <pubDate>Sun, 12 Jul 2026 15:48:58 +0000</pubDate>
      <link>https://dev.to/junhyun-dev/sourcehashro-gateun-ibryeog-jaeceorireul-anjeonhage-skiphagi-4i9b</link>
      <guid>https://dev.to/junhyun-dev/sourcehashro-gateun-ibryeog-jaeceorireul-anjeonhage-skiphagi-4i9b</guid>
      <description>&lt;h1&gt;
  
  
  source_hash로 같은 입력 재처리를 안전하게 skip하기
&lt;/h1&gt;

&lt;p&gt;작은 데이터 파이프라인도 한 번만 실행된다고 가정하면 금방 거짓말이 된다.&lt;/p&gt;

&lt;p&gt;실제로는 같은 파일을 다시 실행할 수 있다. 실패한 run을 재시도할 수도 있고, 과거 날짜를 backfill할 수도 있고, 운영자가 실수로 같은 입력을 다시 넣을 수도 있다. 이때 결과가 중복되면 gold metric은 더 이상 믿을 수 없다.&lt;/p&gt;

&lt;p&gt;이 글은 개인 포트폴리오 프로젝트 &lt;code&gt;manufacturing-data-platform-mini&lt;/code&gt;에서 &lt;code&gt;source_hash&lt;/code&gt;를 이용해 같은 입력 재처리를 안전하게 skip하도록 만든 작은 설계 판단을 정리한 글이다. 데이터는 모두 synthetic이며, production platform이 아니라 검증 가능한 mini slice다.&lt;/p&gt;

&lt;h2&gt;
  
  
  1. Scenario
&lt;/h2&gt;

&lt;p&gt;같은 &lt;code&gt;business_date&lt;/code&gt;의 제조/로봇 이벤트 파일을 다시 처리해야 하는 상황이 있다.&lt;/p&gt;

&lt;p&gt;예:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;retry:
  앞 run이 중간에 실패해서 다시 실행한다.

backfill:
  과거 날짜를 다시 채운다.

operator mistake:
  같은 파일을 실수로 다시 실행한다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;단순히 매번 append하면 같은 날짜의 gold metric이 중복될 수 있다.&lt;/p&gt;

&lt;h2&gt;
  
  
  2. Decision Pressure
&lt;/h2&gt;

&lt;p&gt;단순 CSV pipeline은 보통 이렇게 끝난다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;CSV 읽기 -&amp;gt; silver 만들기 -&amp;gt; gold 집계 -&amp;gt; 결과 저장
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;하지만 운영 관점에서는 질문이 생긴다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;이 입력은 전에 처리한 파일과 같은가?
같은 파일을 다시 돌리면 중복 output이 생기나?
다른 파일로 같은 날짜를 다시 돌리면 어떻게 구분하나?
어떤 run이 어떤 source에서 만들어졌나?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;그래서 재실행을 판단할 identity가 필요했다.&lt;/p&gt;

&lt;h2&gt;
  
  
  3. Options
&lt;/h2&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;option&lt;/th&gt;
&lt;th&gt;result&lt;/th&gt;
&lt;th&gt;problem&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;always append&lt;/td&gt;
&lt;td&gt;모든 run 결과를 계속 추가&lt;/td&gt;
&lt;td&gt;같은 입력 재실행 시 중복&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;always overwrite&lt;/td&gt;
&lt;td&gt;결과를 항상 덮어씀&lt;/td&gt;
&lt;td&gt;이전 결과/원인 추적이 약함&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;skip by business_date only&lt;/td&gt;
&lt;td&gt;같은 날짜면 무조건 skip&lt;/td&gt;
&lt;td&gt;정정 파일을 반영할 수 없음&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;skip by &lt;code&gt;dataset_id + business_date + source_hash&lt;/code&gt;
&lt;/td&gt;
&lt;td&gt;같은 입력만 no-op&lt;/td&gt;
&lt;td&gt;정정 파일은 새 run으로 처리 가능&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;p&gt;이 프로젝트의 Slice1은 마지막 선택지를 쓴다.&lt;/p&gt;

&lt;h2&gt;
  
  
  4. Decision
&lt;/h2&gt;

&lt;p&gt;현재 mini pipeline은 입력 파일의 content hash를 &lt;code&gt;source_hash&lt;/code&gt;로 계산한다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;idempotency key:
  dataset_id + business_date + source_hash
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이미 성공한 run이 있으면 새로 처리하지 않고 기존 run을 재사용한다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;same dataset_id
same business_date
same source_hash
prior successful run exists
=&amp;gt; status = skipped
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이 선택은 작지만 중요하다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;같은 파일 재실행:
  skip -&amp;gt; 중복 없음

같은 날짜의 정정 파일:
  source_hash가 다름 -&amp;gt; skip하지 않고 새 run으로 처리 가능
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;단, 여기서 조심해야 할 경계가 있다. Slice1은 다른 &lt;code&gt;source_hash&lt;/code&gt;를 새 run으로 처리할 수 있지만, 이전 gold partition을 원자적으로 교체하는 Iceberg-style overwrite까지 구현한 것은 아니다. 그 문제는 다음 Slice2의 &lt;code&gt;business_date&lt;/code&gt; partition overwrite 주제다.&lt;/p&gt;

&lt;h2&gt;
  
  
  5. Evidence
&lt;/h2&gt;

&lt;p&gt;관련 코드와 검증 evidence:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;src/manufacturing_data_platform/pipeline/lakehouse.py
tests/test_lakehouse_pipeline.py
VERIFICATION_LOG.md
README.md
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;검증 로그:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;2026-07-08 publication readiness check:
  pytest: 33 passed
  lakehouse JSON CLI: passed, status=processed, quality_passed=true
  EAV JSON CLI: passed, status=processed, quality_passed=true
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이 repo의 테스트는 같은 입력 재실행 시 &lt;code&gt;status="skipped"&lt;/code&gt;가 되는 경로를 구체적으로 확인한다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;tests/test_lakehouse_pipeline.py
  test_rerun_same_source_and_date_is_skipped_mongo
  test_rerun_same_source_and_date_is_skipped_json
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;Mongo backend 쪽 테스트는 재사용 감사 흔적인 &lt;code&gt;reuse_count&lt;/code&gt;가 증가하는지도 확인한다. 즉 단순히 "두 번째 실행을 무시"하는 것이 아니라, 같은 성공 run을 재사용했다는 기록을 남긴다.&lt;/p&gt;

&lt;p&gt;또한 JSON catalog backend로 Mongo 없이도 CLI smoke run을 검증했다.&lt;/p&gt;

&lt;h2&gt;
  
  
  6. Limitations
&lt;/h2&gt;

&lt;p&gt;정직한 한계:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;이건 production lakehouse가 아니다.
Spark/Iceberg는 아직 구현되지 않았다.
pyspark도 현재 repo 환경에 설치되어 있지 않다.
Kafka streaming도 없다.
real Mongo runtime verification은 환경 제약으로 미검증이다.
Airflow runtime trigger도 미검증이다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;현재 claim은 이렇게 제한한다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;synthetic CSV 기반 mini pipeline에서
source_hash를 이용한 idempotent rerun을 구현하고 테스트/CLI로 검증했다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h2&gt;
  
  
  7. Why This Connects To Iceberg Later
&lt;/h2&gt;

&lt;p&gt;Slice1의 skip 전략은 같은 입력 재실행을 안전하게 막는다.&lt;/p&gt;

&lt;p&gt;하지만 다음 질문이 남는다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;정정된 파일로 같은 business_date를 다시 처리해야 한다면?
append하면 중복된다.
skip하면 정정을 반영하지 못한다.
overwrite하면 어디까지 덮어써야 하는가?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이 질문이 Slice2의 Spark/Iceberg 주제로 이어진다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;same source_hash:
  skip

different source_hash for same business_date:
  현재 Slice1에서는 새 run으로 처리 가능
  중복 없는 partition 교체는 Iceberg partition atomic overwrite 후보
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;그래서 다음 글 후보는:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;skip에서 partition overwrite로: business_date 재처리를 Iceberg로 다시 표현하기
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h2&gt;
  
  
  8. 정리
&lt;/h2&gt;

&lt;p&gt;같은 파일을 다시 처리하는 retry와, 같은 날짜의 정정 파일을 다시 처리하는 correction은 다른 문제다.&lt;/p&gt;

&lt;p&gt;이 slice에서는 &lt;code&gt;dataset_id + business_date + source_hash&lt;/code&gt;로 같은 입력을 skip하고, &lt;code&gt;reuse_count&lt;/code&gt;로 재사용 흔적을 남겼다. 정정 파일을 중복 없이 교체하는 문제는 아직 이 글의 구현 범위가 아니며, 다음 Spark/Iceberg slice에서 &lt;code&gt;business_date&lt;/code&gt; partition overwrite로 다룬다.&lt;/p&gt;

</description>
      <category>dataengineering</category>
      <category>python</category>
      <category>etl</category>
      <category>learning</category>
    </item>
    <item>
      <title>gold 숫자가 이상할 때 source_hash, quality, lineage로 원인 좁히기</title>
      <dc:creator>박준현</dc:creator>
      <pubDate>Fri, 10 Jul 2026 06:28:50 +0000</pubDate>
      <link>https://dev.to/junhyun-dev/gold-susjaga-isanghal-ddae-sourcehash-quality-lineagero-weonin-jobhigi-nce</link>
      <guid>https://dev.to/junhyun-dev/gold-susjaga-isanghal-ddae-sourcehash-quality-lineagero-weonin-jobhigi-nce</guid>
      <description>&lt;h1&gt;
  
  
  gold 숫자가 이상할 때 source_hash, quality, lineage로 원인 좁히기
&lt;/h1&gt;

&lt;p&gt;데이터 플랫폼에서 중요한 질문은 "gold 숫자를 만들었는가"에서 끝나지 않는다.&lt;/p&gt;

&lt;p&gt;운영 중에는 오히려 이런 질문이 더 자주 나온다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;이 defect_rate가 왜 이렇게 높지?
이 gold row는 어느 source에서 왔지?
row가 처리 중 사라진 건가, 정상 필터링인가?
schema가 바뀐 상태로 계산된 건가?
quality check는 통과했나?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;&lt;code&gt;manufacturing-data-platform-mini&lt;/code&gt;의 B4 slice는 이 질문을 작게 다룬다. 새 Spark/Iceberg 엔진을 붙이지 않고, 이미 남겨둔 JSON catalog/lineage/quality evidence를 read-only operator report로 읽는다.&lt;/p&gt;

&lt;p&gt;이 report는 이상치를 자동 탐지하지 않는다. 숫자를 해석할 provenance/quality 맥락을 제공해서 사람이 원인 후보를 좁히게 하는 도구다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Scenario
&lt;/h2&gt;

&lt;p&gt;분석가가 &lt;code&gt;business_date=2026-06-29&lt;/code&gt;의 &lt;code&gt;defect_rate&lt;/code&gt;가 이상하다고 말한다.&lt;/p&gt;

&lt;p&gt;운영자는 raw CSV를 바로 열기 전에, 먼저 metadata와 evidence로 원인을 좁히고 싶다.&lt;/p&gt;

&lt;p&gt;확인 순서는 이렇다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;1. gold row grain 확인
2. 해당 business_date의 successful run 확인
3. run_id / source_hash / schema_hash 확인
4. quality checks에서 fail/warn 확인
5. lineage path로 gold -&amp;gt; silver -&amp;gt; bronze -&amp;gt; source 역추적
6. row count / conservation / schema drift를 보고 원인 후보 좁히기
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h2&gt;
  
  
  Decision Pressure
&lt;/h2&gt;

&lt;p&gt;단순 CSV 변환 스크립트는 보통 output만 남긴다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;input.csv -&amp;gt; gold.csv
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;하지만 운영 질문은 output만으로 답하기 어렵다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;이 gold.csv는 어떤 입력에서 왔나?
동일 source를 재실행한 건가, 다른 source로 새로 처리한 건가?
source row 5개가 silver 3개가 된 이유는 정상 필터링/중복 제거인가?
집계가 units/defects를 보존했나?
schema drift가 있었나?
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;그래서 이 프로젝트는 transform output과 함께 run evidence를 남긴다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;source_hash
schema_hash
quality checks
stats
layer parent links
catalog state
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이 문제는 data lineage 도구들이 다루는 전형적인 운영 질문의 축소판이다. 예를 들어 &lt;a href="https://openlineage.io/" rel="noopener noreferrer"&gt;OpenLineage&lt;/a&gt;는 dataset/job/run metadata를 추적해 문제의 원인과 변경 영향 이해를 돕는 표준을 제공하고, &lt;a href="https://docs.databricks.com/aws/en/data-governance/unity-catalog/data-lineage" rel="noopener noreferrer"&gt;Databricks Unity Catalog lineage&lt;/a&gt;는 downstream 결과가 이상할 때 upstream source를 추적하는 root-cause investigation을 lineage의 사용 사례로 설명한다.&lt;/p&gt;

&lt;p&gt;이 프로젝트는 그런 시스템을 구현한 것이 아니라, 같은 문제를 path-level JSON evidence로 작게 연습한다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Options
&lt;/h2&gt;

&lt;div class="table-wrapper-paragraph"&gt;&lt;table&gt;
&lt;thead&gt;
&lt;tr&gt;
&lt;th&gt;Option&lt;/th&gt;
&lt;th&gt;장점&lt;/th&gt;
&lt;th&gt;문제&lt;/th&gt;
&lt;th&gt;판단&lt;/th&gt;
&lt;/tr&gt;
&lt;/thead&gt;
&lt;tbody&gt;
&lt;tr&gt;
&lt;td&gt;raw CSV 직접 열기&lt;/td&gt;
&lt;td&gt;가장 직접적&lt;/td&gt;
&lt;td&gt;매번 수동, run/source identity를 놓치기 쉬움&lt;/td&gt;
&lt;td&gt;보조 수단&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;gold 파일만 보기&lt;/td&gt;
&lt;td&gt;빠름&lt;/td&gt;
&lt;td&gt;원인 추적 불가&lt;/td&gt;
&lt;td&gt;부족&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;full lineage system&lt;/td&gt;
&lt;td&gt;강력함&lt;/td&gt;
&lt;td&gt;v0에 과함, OpenLineage backend 미구현&lt;/td&gt;
&lt;td&gt;backlog&lt;/td&gt;
&lt;/tr&gt;
&lt;tr&gt;
&lt;td&gt;read-only operator report&lt;/td&gt;
&lt;td&gt;작고 검증 가능, 기존 evidence 재사용&lt;/td&gt;
&lt;td&gt;path-level까지만 가능&lt;/td&gt;
&lt;td&gt;선택&lt;/td&gt;
&lt;/tr&gt;
&lt;/tbody&gt;
&lt;/table&gt;&lt;/div&gt;

&lt;h2&gt;
  
  
  Decision
&lt;/h2&gt;

&lt;p&gt;이 slice에서는 read-only operator report를 선택했다.&lt;/p&gt;

&lt;p&gt;명령은 작다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight shell"&gt;&lt;code&gt;&lt;span class="nv"&gt;PYTHONPATH&lt;/span&gt;&lt;span class="o"&gt;=&lt;/span&gt;src python &lt;span class="nt"&gt;-m&lt;/span&gt; manufacturing_data_platform.pipeline.operator_report &lt;span class="se"&gt;\&lt;/span&gt;
  &lt;span class="nt"&gt;--output-dir&lt;/span&gt; /tmp/manufacturing-mini-operator-report-cli &lt;span class="se"&gt;\&lt;/span&gt;
  &lt;span class="nt"&gt;--business-date&lt;/span&gt; 2026-06-29
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이 report는 JSON catalog state를 읽어서 아래를 보여준다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;gold grain
run_id
source_hash
schema_hash
quality summary
row counts
lineage trace: gold -&amp;gt; silver -&amp;gt; bronze -&amp;gt; source
claim boundary
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;또한 report 자체가 자기 한계를 같이 출력한다. 블로그에서만 정직한 것이 아니라, 산출물도 claim boundary를 포함한다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Evidence
&lt;/h2&gt;

&lt;p&gt;구현 evidence:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;src/manufacturing_data_platform/pipeline/operator_report.py
tests/test_operator_report.py
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;검증 로그:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;2026-07-10 — Operator evidence report slice
pytest: 35 passed
lakehouse JSON CLI: passed
operator evidence report CLI: passed
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;실제 출력 일부:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight json"&gt;&lt;code&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"gold_grain"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"dataset_id"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"manufacturing_daily_metrics"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"row_grain"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="s2"&gt;"business_date"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="s2"&gt;"plant_id"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="s2"&gt;"line_id"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="s2"&gt;"product_code"&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;],&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"metrics"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="s2"&gt;"units_produced"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="s2"&gt;"defect_count"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="s2"&gt;"defect_rate"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="s2"&gt;"avg_cycle_time_ms"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="s2"&gt;"closing_status"&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;]&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="p"&gt;},&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"run"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"run_id"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"2026-06-29-20260710T033849Z-73005763"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"source_hash"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"b3ffc4acdd909db1b5a87db7155f76589dfdd4cedb4b00cedf18506f12948604"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"schema_hash"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"4414a80b70b7b386a13ee705a33aa9c99bd29cb3a717bba9d2c5f0bc892d3126"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"quality_passed"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="kc"&gt;true&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"reuse_count"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="p"&gt;},&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"quality_summary"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"total_checks"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;8&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"pass_count"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;8&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"warn_count"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"fail_count"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;0&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"failed_checks"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;[],&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"warning_checks"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;[],&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"rca_focus_checks"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="nl"&gt;"name"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"row_count_source_to_silver"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="nl"&gt;"status"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"pass"&lt;/span&gt;&lt;span class="p"&gt;},&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="nl"&gt;"name"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"unit_conservation_silver_to_gold"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="nl"&gt;"status"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"pass"&lt;/span&gt;&lt;span class="p"&gt;},&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="nl"&gt;"name"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"freshness_business_date"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="nl"&gt;"status"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"pass"&lt;/span&gt;&lt;span class="p"&gt;},&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="nl"&gt;"name"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"schema_drift"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="nl"&gt;"status"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"pass"&lt;/span&gt;&lt;span class="p"&gt;}&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;]&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="p"&gt;},&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"lineage_trace"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="nl"&gt;"name"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"gold"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="nl"&gt;"parents"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="s2"&gt;".../silver/manufacturing_events.csv"&lt;/span&gt;&lt;span class="p"&gt;]},&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="nl"&gt;"name"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"silver"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="nl"&gt;"parents"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="s2"&gt;".../bronze/manufacturing_events.csv"&lt;/span&gt;&lt;span class="p"&gt;]},&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="nl"&gt;"name"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"bronze"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="nl"&gt;"parents"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="s2"&gt;"data/raw/manufacturing_events.csv"&lt;/span&gt;&lt;span class="p"&gt;]},&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"name"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"source"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"path"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"data/raw/manufacturing_events.csv"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"hash"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="s2"&gt;"b3ffc4acdd909db1b5a87db7155f76589dfdd4cedb4b00cedf18506f12948604"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="nl"&gt;"row_count"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="mi"&gt;5&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;}&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="p"&gt;],&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="nl"&gt;"claim_boundary"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;{&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"supports"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="s2"&gt;"table/path-level lineage"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="s2"&gt;"operator-inspectable run evidence"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="s2"&gt;"source/schema identity trace"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="s2"&gt;"quality check summary"&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;],&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="nl"&gt;"does_not_support"&lt;/span&gt;&lt;span class="p"&gt;:&lt;/span&gt;&lt;span class="w"&gt; &lt;/span&gt;&lt;span class="p"&gt;[&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="s2"&gt;"column-level lineage"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="s2"&gt;"OpenLineage backend integration"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="s2"&gt;"interactive lineage UI"&lt;/span&gt;&lt;span class="p"&gt;,&lt;/span&gt;&lt;span class="w"&gt;
      &lt;/span&gt;&lt;span class="s2"&gt;"production incident workflow"&lt;/span&gt;&lt;span class="w"&gt;
    &lt;/span&gt;&lt;span class="p"&gt;]&lt;/span&gt;&lt;span class="w"&gt;
  &lt;/span&gt;&lt;span class="p"&gt;}&lt;/span&gt;&lt;span class="w"&gt;
&lt;/span&gt;&lt;span class="p"&gt;}&lt;/span&gt;&lt;span class="w"&gt;
&lt;/span&gt;&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이 출력을 시나리오에 대입하면:&lt;/p&gt;

&lt;ul&gt;
&lt;li&gt;
&lt;code&gt;quality_summary.fail_count=0&lt;/code&gt;, &lt;code&gt;warn_count=0&lt;/code&gt; -&amp;gt; 유실, 보존 불일치, schema drift warning이 없었다. 즉 이상한 &lt;code&gt;defect_rate&lt;/code&gt;는 파이프라인 버그보다 실제 입력 데이터에서 온 값일 가능성이 크다.&lt;/li&gt;
&lt;li&gt;
&lt;code&gt;lineage_trace&lt;/code&gt;의 &lt;code&gt;source.hash&lt;/code&gt;와 &lt;code&gt;row_count=5&lt;/code&gt; -&amp;gt; 어느 파일의 5개 row에서 왔는지 raw를 열기 전에 고정할 수 있다.&lt;/li&gt;
&lt;li&gt;
&lt;code&gt;run.reuse_count=0&lt;/code&gt; -&amp;gt; 이전 successful run을 재사용한 skip이 아니라 새로 처리된 run이다.&lt;/li&gt;
&lt;/ul&gt;

&lt;p&gt;report 하나로 "숫자가 이상하다"를 "파이프라인 check는 정상이고, 출처는 이 source_hash의 5-row file"까지 좁힌다.&lt;/p&gt;

&lt;h2&gt;
  
  
  Limitations
&lt;/h2&gt;

&lt;p&gt;이건 production lineage system이 아니다.&lt;/p&gt;

&lt;p&gt;명확한 한계:&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;이상치 자동 탐지 아님
자동 RCA 아님
latest successful run evidence 조회용 — failure-state forensics는 backlog
column-level lineage 아님
OpenLineage backend 통합 아님
interactive lineage UI 아님
production incident workflow 아님
real Mongo runtime 검증 아님
Airflow runtime trigger 검증 아님
Spark/Iceberg 구현 아님
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;p&gt;이 slice의 목적은 작다.&lt;br&gt;
&lt;/p&gt;

&lt;div class="highlight js-code-highlight"&gt;
&lt;pre class="highlight plaintext"&gt;&lt;code&gt;기존 run/catalog/quality/lineage evidence를 operator가 읽을 수 있는 형태로 묶는다.
&lt;/code&gt;&lt;/pre&gt;

&lt;/div&gt;



&lt;h2&gt;
  
  
  정리
&lt;/h2&gt;

&lt;p&gt;output만 남기면 "왜 이 숫자지?"에 답하기 어렵다. output과 함께 &lt;code&gt;source_hash&lt;/code&gt;, quality result, row counts, lineage evidence를 남기면 raw 파일을 열기 전에 "어디서 왔고, 파이프라인 check는 정상이었는지"까지 원인 후보를 좁힐 수 있다.&lt;/p&gt;

&lt;p&gt;이 report는 anomaly detection이나 자동 RCA가 아니다. 의심스러운 지표를 run, source, quality check, path-level lineage 맥락으로 좁히는 read-only evidence view다.&lt;/p&gt;

&lt;p&gt;코드: &lt;a href="https://github.com/junhyun-dev/manufacturing-data-platform-mini" rel="noopener noreferrer"&gt;github.com/junhyun-dev/manufacturing-data-platform-mini&lt;/a&gt;&lt;/p&gt;

</description>
      <category>dataengineering</category>
      <category>python</category>
      <category>etl</category>
      <category>learning</category>
    </item>
  </channel>
</rss>
