You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
AKS 위 Airflow DAG가 Azure DocumentDB(vCore) change stream을 주기적으로 읽어 ADLS Gen2에 Parquet으로 쓸 때 이벤트 누락과 중복이 없는지 실제 구독에서 확인한 topic을 추가합니다. docs/services/azure-documentdb/change-streams/ 아래 entry 1개, 자식 문서 1개, runnable sample 1개로 구성했습니다.
경로
유형
내용
index.md
research
change stream 소개, 확인한 구성(구성 다이어그램 SVG), 동작 여부와 지켜야 할 조건, 공식 예제가 놓친 부분, Learn과 다르게 동작한 부분, 운영 전에 더 확인할 항목
measurements/index.md
research
10분 구간 지연, 크기 한도별(50·80·100·200 MB) 파일 크기와 메모리, 80 MB Airflow 실행과 재시도, 이전 5분 주기 구성의 지연, 업로드 직후 장애와 재시도, 400 MB 활성 change log를 넘긴 재개, Parquet 파일 크기, 상시 PyMongo consumer 기준선, 옵션별 동작의 측정 수치
이벤트를 wallTime 기준 10분 구간(UTC)으로 묶고 압축 전 행 크기가 50 MB를 넘기 전에 파일을 나눕니다. 열린 구간은 쓰지 않고 DAG는 구간 종료 2분 뒤 CronTriggerTimetable로 실행합니다. 자르는 위치가 이벤트에만 의존하므로 재시도가 같은 파일을 덮어씁니다.
새로 배포한 환경에서 지속 부하(1,794,000건), 업로드 직후 장애 재시도, 열린 구간 보류, 메모리를 측정했고 모두 중복·누락 0이었습니다. 백로그 재개와 consumer 기준선은 이전 구성의 측정으로 표시해 남겼습니다. ADLS 모범 사례를 출처에 추가했고 Airflow on AKS, Fabric open mirroring, pandas·pyarrow ADLS 쓰기는 "더 읽을 문서"에 링크만 두었습니다.
리뷰를 반영해 측정 이후 바꾼 sample 코드(첫 실행 시작 위치 저장, lock 파일 lease, dt=unknown 파티션, 생성기 실패 처리, 검증기 순서 역전 판정, 사용자 할당 AKS ID, Airflow RBAC 축소 등)는 Azure에 다시 배포해 확인했습니다. README 기본값으로 provision, r1, 첫 DAG 실행, 100,000개 적재, 업로드 직후 장애 재시도, lease 경쟁을 실행했고 모두 기대대로 동작했습니다. wallTime이 없는 이벤트와 순서 역전 경로는 이 클러스터에서 생기지 않아 로컬 fake로만 확인했습니다.
50 MB에서 자른 파일이 29 MB 안팎에 그쳐 한도를 50·100·200 MB로 바꿔 같은 이벤트를 내보냈습니다. 파일 크기는 한도에 비례했고(1 KB 채움 약 0.58배, 채움 없음 약 0.17배) 처리 시간은 한도와 관계없이 38–41초였습니다. 1 KB 채움 문서에서도 파일이 50 MB를 넘지 않는 80 MB를 기본값으로 정했습니다.
압축 전 한도
1 KB 채움 파일
최대 RSS
채움 없음 파일
최대 RSS
50 MB
28.1–29.3 MB
305 MiB
8.3–8.5 MB
253 MiB
80 MB
43.7–46.7 MB
336 MiB
13.3–13.4 MB
재지 않음
100 MB
57.6–58.4 MB
390 MiB
16.6–16.8 MB
335 MiB
200 MB
115.4–116.7 MB
620 MiB
32.8 MB
519 MiB
행을 10,000건마다 Arrow 배치로 옮기고 업로드할 때 버퍼 사본을 만들지 않게 바꿔 50 MB 한도의 최대 RSS가 358/386 MiB에서 305/253 MiB로 줄었습니다.
80 MB 이미지로 Airflow를 갱신한 뒤 밀린 세 구간 2,185,000건을 정기 실행 한 번이 175.8초에 파일 22개로 썼고 업로드 직후 장애 재시도도 690,000건 중복·누락 0이었습니다.
재시도가 첫 시도와 다른 위치에서 잘라도 중복이 생기지 않는 것을 로컬 fake로 확인해 문서의 지나친 주장("압축 후 크기로 자르면 중복")을 고쳤습니다.
DAG를 멈췄다 켜면 놓친 정기 실행이 수동 실행보다 먼저 돌아 README의 재시도 시험 절차를 DAG를 켠 채 실행하는 방식으로 바꿨습니다.
압축이 거의 안 되는 데이터는 파일이 80 MB에 가까워질 수 있어 50 MB 아래를 지켜야 하면 CS_CHUNK_BYTES=50000000을 쓰도록 적었습니다.
함께 변경한 파일
docs-taxonomy.yml: 서비스 azure-documentdb, 기술 airflow, mongodb 추가
공개 문서 체크리스트
사례·가이드·실습·리서치 중 올바른 문서 유형과 폴더를 선택했습니다.
필수 front matter와 docs-taxonomy.yml의 분류 값을 사용했습니다.
Microsoft/Azure 기술 주장을 Microsoft Learn MCP와 verify-with-microsoft-learn skill로 원문까지 확인했습니다.
실제 구독·테넌트 ID, 비밀, 내부 URL·IP·호스트명 등 공개 안전성 문제를 제거했습니다.
필요한 경우 현재 가이드와 시점 고정 사례를 분리하고 서로 연결했습니다.
아래 로컬 검증을 통과했습니다.
리소스 이름은 문서에 쓰지 않고 sample은 azd 환경 값으로 받습니다. maxAwaitTimeMS 동작처럼 Learn에 설명이 없는 부분은 이 클러스터의 관측 결과로 표시했습니다.
- 프라이빗 엔드포인트 뒤 Azure DocumentDB(M30)에서 PyMongo change stream의
지원 범위, 재개, 중복과 지연을 AKS consumer로 측정한 lab 추가
- AKS 위 Airflow DAG가 5분마다 change stream을 읽어 ADLS Gen2에 Parquet으로
쓰는 구성의 지연, 재시도 중복, 400 MB 활성 change log를 넘긴 재개와 파일
크기를 측정한 research 문서와 구성 다이어그램 추가
- maxAwaitTimeMS가 getMore 실행 제한으로 적용되어 큰 백로그 재개가 code 50으로
반복 실패하는 동작을 기록하고 내보내기 코드의 기본값에서 제외
- runnable sample(python-aks): Bicep, consumer·생성기·검증기, 내보내기와 DAG
- taxonomy에 azure-documentdb 서비스와 airflow, mongodb 기술 추가
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
- azure.yaml과 구독 범위 main.bicep을 추가하고 기존 템플릿은 cluster.bicep으로 옮김
- 관리자 비밀번호는 azd .env에 넣지 않고 셸 환경 변수로만 전달
- 이미지, 파드와 Airflow는 README의 az acr build, kubectl, helm 절차로 배포
- README 절차대로 새 환경에서 r1과 af1을 다시 실행해 누락·중복 0건 확인
- 수동으로 만드는 namespace와 서비스 계정을 매니페스트에서 제거
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
- 주제 첫 문서를 research로 바꾸고 기능 소개, 동작 여부, 공식 예제가 놓친 부분,
운영 전 확인 항목 위주로 정리
- 측정 수치와 상시 consumer 기준선, 옵션별 관찰은 measurements 하위 문서로 이동
- 배포와 실행 절차는 sample README로 옮기고 측정 시나리오 매개변수를 추가
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
- 첫 실행이 읽기 전에 시작 시각을 저장해 재시도가 같은 위치에서 읽도록 함
- 내보내기 실행이 lock 파일 lease를 잡아 겹친 실행의 파일 덮어쓰기를 막음
- wallTime이 없는 이벤트는 dt=unknown 파티션에 써서 재시도 경로를 고정
- generator가 worker 실패를 failed로 기록하고 exit 1로 종료
- verify.py, verify_lake.py가 문서별 순서 역전을 실패로 처리
- README의 r1, r4 생성기 매개변수를 측정값에 맞춤
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
시스템 할당 ID는 클러스터 생성 후에야 생겨 Network Contributor 역할이
클러스터보다 늦게 붙었다. 사용자 할당 ID를 먼저 만들고 snet-aks 범위로
역할을 부여한 뒤 AKS가 그 역할 부여에 의존하도록 바꾼다.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
Keying observations by operationType overwrites the ordinary update event when this cluster emits the following replacement as another update, which is the behavior documented by this PR. The probe output therefore cannot retain evidence for both events. Preserve all four observations in a list (including each operation type) instead of collapsing them into this dictionary.
Unlike the surrounding probes, this dereferences the event after a timed wait without handling None. If the cluster returns no update event within WAIT_S, the check records a Python TypeError instead of the intended observation, losing whether the fields were absent. Guard the timeout result as the other checks do.
Escape RUN_ID before constructing the verification regex
RUN_ID is inserted into a regular expression without escaping, although the generator accepts arbitrary non-empty IDs and the lake verifier treats them literally. An ID such as r.1 also matches rX1, so rows from another run can produce false unexpected-event failures. Escape the ID before constructing the anchored regex.
Align sample node count with documented test environment
The documented provisioning commands do not override this value, so the runnable sample creates two nodes, while the published test environment and architecture state four Standard_D4s_v6 nodes (index.md:161, architecture.svg:12). That prevents readers from reproducing the reported performance with the documented procedure; either use four here or explicitly document how the measured deployment overrode the sample default.
Align document classification and path with PR description
The PR description classifies the entry as a lab, but this front matter publishes it as research (and the described airflow-parquet/ child is actually measurements/). Align the PR description or the document classification/path so reviewers and lifecycle automation have one authoritative scope.
- Airflow scheduler Role에서 쓰지 않는 pods/exec 권한 제거
- probe의 이벤트 형태 확인이 같은 operationType 이벤트를 덮어쓰지 않도록 목록으로 기록
- probe의 기본 update 확인이 이벤트를 받지 못하면 FAIL로 기록
- verify.py가 RUN_ID를 정규식에 넣기 전에 escape
- AKS 노드 수 기본값을 측정 환경과 같은 4개로 변경
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
hellices
changed the title
docs(azure-documentdb): change stream 실습과 Airflow Parquet 적재 측정 추가
docs(azure-documentdb): change stream을 Airflow DAG로 Parquet에 적재하는 검증 추가
Oct 2, 2026
리뷰 반영 후 sample을 Azure에 다시 배포해 첫 실행 시작 위치, stream
lease 경쟁, 100,000개 적재와 업로드 직후 장애 재시도를 실행했다.
로컬 fake로만 확인했다는 서술을 실측 결과로 바꾸고 wallTime 없는
이벤트와 순서 역전 경로만 로컬 확인으로 남긴다. pause 상태에서
trigger한 run은 unpause해야 실행되므로 retry 절차를 고친다.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
앞선 코멘트와 RBAC 스레드 답글에서 "Azure에서 다시 실행하지 못했다"고 적었던 부분을 실제 배포로 확인했습니다. 22270d3 sample을 README 기본값(eastus2, 노드 4개)으로 새 리소스 그룹에 배포했습니다.
항목
결과
azd provision
13분 37초에 성공. 사용자 할당 ID에 서브넷 권한을 먼저 준 구성으로 AKS 생성, 노드 4개 Ready
r1 (consumer, 10,000개)
23,000/23,000건. 누락·중복·순서 역전 0
probe event_shapes
PASS. 이벤트 4건이 목록으로 남음(replace는 update로 도착)
verify.py 접두사 조회
점이 들어간 RUN_ID v.1로 2,300/2,300건 대조
pods/exec 없는 RBAC
DAG 실행 7회 모두 성공, 권한 거부 없음
첫 DAG 실행
checkpoint 없이 시작 위치를 저장하고 0건으로 종료. 다음 실행부터 이어 읽음
af1 (100,000개, 정기 실행)
Parquet 230,000/230,000건. 중복·누락·순서 역전 0
af2 업로드 직후 장애 재시도
첫 시도가 청크 1개 업로드 후 종료. 재시도가 같은 파일을 덮어쓰고 230,000/230,000건, 중복 0
lease 경쟁
다른 파드가 lease를 잡은 동안 띄운 export가 91초 기다린 뒤 LeaseAlreadyPresent로 exit 1. change stream은 열지 않음
확인하지 못한 경로도 있습니다. wallTime이 없는 이벤트의 dt=unknown 파티션과 검증기의 순서 역전 실패 처리는 이 클러스터에서 생기지 않아 여전히 로컬 fake로만 확인했습니다.
실행 중 README 절차 하나가 틀린 것을 발견했습니다. pause 상태에서 trigger한 DAG run은 queued로 남고 unpause한 뒤에야 다음 정기 실행보다 먼저 실행됩니다. 6b0ab49에서 retry 절차를 이 순서로 고쳤고 index.md와 README의 "로컬 fake로만 확인" 서술도 위 결과로 바꿨습니다. 시험 리소스는 확인 후 삭제합니다.
- 이벤트를 wallTime 기준 10분 구간(UTC)으로 묶고 압축 전 50 MB 전에
파일을 나눈다. 열린 구간은 쓰지 않고 다음 실행으로 넘긴다.
- DAG를 구간 종료 2분 뒤 CronTriggerTimetable로 실행한다.
- 열린 구간 이벤트를 읽은 실행은 idle checkpoint를 저장하지 않는다.
- verify_lake.py가 구간이 섞인 파일과 한도를 넘은 파일을 실패로 본다.
- 내보내기 파드 메모리를 요청 512Mi, 제한 1Gi로 낮춘다.
- 새 배포에서 측정한 지연, 재시도, 열린 구간 보류, 파일 크기와 메모리를
문서에 반영하고 ADLS 모범 사례를 출처에 추가한다.
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This generalization includes deletes, but the preceding measurement says the 8,713 delete rows had no document body. Limit the statement to the observed insert and update events so it does not imply that every change event carries fullDocument.
The probe inserts a document with one 7 MiB field and only reaches about 14 MiB after the update adds the second field, so the insert event was not for a 14 MiB document. Distinguish the two sizes to keep the reported observation accurate.
airflow dags list-runs requires the DAG ID through -d/--dag-id; passing it positionally causes this verification command to fail with an unrecognized argument. Use the documented named option.
- Chunk가 wallTime 없는 이벤트와 날짜가 있는 구간을 한 파일에 섞지 않도록 구간 값을 그대로 비교
- verify_lake.py가 wallTime 없는 행도 하나의 구간으로 세어 섞인 파일을 실패로 판정
- 측정 문서의 fullDocument 설명을 insert·update로 한정하고 큰 문서 크기를 7 MiB insert, 약 14 MiB update로 정정
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
measurements/index.md fullDocument 설명: 7c77ef6에서 insert와 update 이벤트에만 문서 전체가 실리고 delete에는 본문이 없다고 고쳤습니다.
measurements/index.md 큰 문서 행: probe.py는 7 MiB 문서를 insert하고 같은 크기 필드를 $set해 약 14 MiB로 키웁니다. "7 MiB 문서의 insert와 약 14 MiB로 커진 update 이벤트 모두 전달됨"으로 정정했습니다.
README airflow dags list-runs change_stream_to_parquet: 바꾸지 않았습니다. Airflow 3.2.2 CLI 참조에서 list-runs의 dag_id는 위치 인자이고 -d/--dag-id 옵션은 없습니다. 이 sample의 Airflow 3.2.2 배포에서도 위치 인자 형태로 실행됐습니다.
- 행을 10,000건마다 Arrow RecordBatch로 옮기고 업로드할 때 버퍼 사본을
만들지 않도록 바꿔 50 MB 한도의 최대 RSS를 358/386 MiB에서 305/253 MiB로 줄임
- CHUNK_BYTES 기본값을 80,000,000으로 올림. 1 KB 채움 문서 파일은
43.7-46.7 MB, 채움 없는 문서는 13.3-13.4 MB
- 50/80/100/200 MB 한도 비교, 80 MB Airflow 백로그 실행과 업로드 직후
장애 재시도 결과를 측정 상세에 추가
- 재시도가 다른 위치에서 잘라도 중복이 생기지 않는다는 점을 반영해
파일 이름 조건 설명을 고침
- Airflow 재시도 시험 절차를 DAG를 켠 채 수동 실행하도록 수정
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
The repository contract at docs/contributing/index.md:268-276 permits verified only after full-source semantic verification. The required Microsoft Learn MCP server is unavailable in this review, so the material Azure claims cannot be reverified through the mandated workflow; keep this as needs-review until that verification is completed.
Keep measurement-page claims as needs-review pending verification
The repository contract at docs/contributing/index.md:268-276 permits verified only after full-source semantic verification. The required Microsoft Learn MCP server is unavailable in this review, so the product claims in this measurement page cannot be reverified through the mandated workflow; keep this as needs-review until that verification is completed.
Chart 1.22.0 requires Helm 3.19.0, but the prerequisites accept any Helm version. Users with an older Helm client will fail at the pinned install command in step 5, so state the minimum version here.
- 청크 경로를 checkpoint에 먼저 기록하고 청크 파일 lease를 잡은 뒤 lease ID로 업로드
- stale_writer_test.py 추가, Azure에서 이전/현재 sample 결과를 문서에 반영
- 메모리 측정 파드와 Airflow DAG 파드의 리소스 조건을 구분
- README 사전 요구 사항에 Helm 3.19.0 이상 명시
Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
verification_status: verified 유지: 상위 문서의 Microsoft 동작 주장은 2026-10-03에 Microsoft Learn MCP로 원문 전체를 다시 조회해 대조했습니다. 청크 파일 lease를 추가하면서 412 오류 코드(LeaseIdMissing, LeaseIdMismatch, LeaseNotPresent)와 lease 기간의 근거로 Path - Update를 두 문서의 official_sources에 더했습니다. lease를 끊을 수 있는 조건은 기존 Lease Blob으로 확인했습니다. 측정 상세의 sources_checked_at도 2026-10-03으로 맞췄습니다.
Helm 버전: Airflow Helm chart 1.22.0 upstream 문서의 요구 사항(Kubernetes v1.30.13+, Helm v3.19.0+)을 확인해 README 사전 요구 사항에 "Helm 3.19.0 or later (Airflow chart 1.22.0 requires it)"를 적었습니다.
검증 범위: 문서 검증 8종(pytest tests/docs, metadata, sources, links, public safety, mkdocs build --strict, search index, pre-pages audit)이 모두 통과했습니다.
CS_CHUNK_BYTES=50000000 caps only the estimated pre-compression row payload; it does not account for Parquet metadata or encoding/compression overhead. Therefore this value does not guarantee the stated strict 50 MB file ceiling. Recommend a measured margin or post-serialization size enforcement instead.
Estimated row size does not guarantee a sub-50 MB Parquet file
CHUNK_BYTES limits the custom row_bytes() estimate, not the serialized Parquet size. That estimate excludes Parquet metadata/encoding overhead, and Snappy can add overhead for incompressible input, so setting it to exactly 50,000,000 cannot guarantee a file below 50 MB. Describe the measured result as such and require a safety margin (or post-serialization enforcement) for a hard limit.
This branch has not been deployed
No deployments
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
변경 내용
AKS 위 Airflow DAG가 Azure DocumentDB(vCore) change stream을 주기적으로 읽어 ADLS Gen2에 Parquet으로 쓸 때 이벤트 누락과 중복이 없는지 실제 구독에서 확인한 topic을 추가합니다.
docs/services/azure-documentdb/change-streams/아래 entry 1개, 자식 문서 1개, runnable sample 1개로 구성했습니다.index.mdmeasurements/index.mdsamples/python-aks/주요 결과
maxAwaitTimeMS=1000이면 이 클러스터가 이를getMore실행 제한으로 적용해 같은 위치에서 code 50으로 반복 실패. 지정하지 않으면 1,380,000건을 누락·중복 없이 저장10분 구간·50 MB 규칙 (69966e1)
이벤트를
wallTime기준 10분 구간(UTC)으로 묶고 압축 전 행 크기가 50 MB를 넘기 전에 파일을 나눕니다. 열린 구간은 쓰지 않고 DAG는 구간 종료 2분 뒤CronTriggerTimetable로 실행합니다. 자르는 위치가 이벤트에만 의존하므로 재시도가 같은 파일을 덮어씁니다.새로 배포한 환경에서 지속 부하(1,794,000건), 업로드 직후 장애 재시도, 열린 구간 보류, 메모리를 측정했고 모두 중복·누락 0이었습니다. 백로그 재개와 consumer 기준선은 이전 구성의 측정으로 표시해 남겼습니다. ADLS 모범 사례를 출처에 추가했고 Airflow on AKS, Fabric open mirroring, pandas·pyarrow ADLS 쓰기는 "더 읽을 문서"에 링크만 두었습니다.
리뷰를 반영해 측정 이후 바꾼 sample 코드(첫 실행 시작 위치 저장, lock 파일 lease,
dt=unknown파티션, 생성기 실패 처리, 검증기 순서 역전 판정, 사용자 할당 AKS ID, Airflow RBAC 축소 등)는 Azure에 다시 배포해 확인했습니다. README 기본값으로 provision, r1, 첫 DAG 실행, 100,000개 적재, 업로드 직후 장애 재시도, lease 경쟁을 실행했고 모두 기대대로 동작했습니다.wallTime이 없는 이벤트와 순서 역전 경로는 이 클러스터에서 생기지 않아 로컬 fake로만 확인했습니다.청크 한도 80 MB (92f3796)
50 MB에서 자른 파일이 29 MB 안팎에 그쳐 한도를 50·100·200 MB로 바꿔 같은 이벤트를 내보냈습니다. 파일 크기는 한도에 비례했고(1 KB 채움 약 0.58배, 채움 없음 약 0.17배) 처리 시간은 한도와 관계없이 38–41초였습니다. 1 KB 채움 문서에서도 파일이 50 MB를 넘지 않는 80 MB를 기본값으로 정했습니다.
CS_CHUNK_BYTES=50000000을 쓰도록 적었습니다.함께 변경한 파일
docs-taxonomy.yml: 서비스azure-documentdb, 기술airflow,mongodb추가공개 문서 체크리스트
docs-taxonomy.yml의 분류 값을 사용했습니다.verify-with-microsoft-learnskill로 원문까지 확인했습니다.리소스 이름은 문서에 쓰지 않고 sample은 azd 환경 값으로 받습니다.
maxAwaitTimeMS동작처럼 Learn에 설명이 없는 부분은 이 클러스터의 관측 결과로 표시했습니다.검증
저장소 루트에서 필수 검증을 모두 실행했습니다(92f3796 기준).
pytest tests/docs1,883건 통과validate_metadata·validate_sources·validate_links·validate_public_safety통과 (76개 문서)mkdocs build --strict통과validate_search_index통과 (76개 문서, 13개 태그)audit_pre_pages통과: 기준선 문서 62/62 보존🤖 Generated with Claude Code