데이터 코파일럿 구축(Build a data copilot)
데이터 코파일럿 구축(Build a data copilot)
이 튜토리얼은 데이터 엔지니어링 코파일럿을 구축합니다: 각 기능이 Cortex Code Agent SDK로 구동되는 작은 웹 앱입니다. SQL을 직접 작성하는 대신 자연어로 원하는 것을 설명하면, 에이전트가 내장 SQL 도구를 사용해 Snowflake에 대한 쿼리를 작성·실행하고, 애플리케이션이 직접 사용할 수 있는 타입화된 구조화 결과를 반환합니다.
초점은 웹 프레임워크가 아니라 SDK 사용입니다. 예시는 브라우저에서 에이전트를 쉽게 구동하므로 Next.js 앱을 사용하지만, 여기의 모든 SDK 패턴은 일반 스크립트, 예약 작업 또는 CI/CD(지속적 통합 및 지속적 배포) 단계에서도 동일하게 작동합니다.
본문
무엇을 구축하나
각각 SDK가 지원하는 세 가지 도구가 있는 데이터 엔지니어링 코파일럿:
- 파이프라인 상태 감사(Pipeline health audit): 스키마를 설명하면 에이전트가 테이블을 발견하고, 행 수를 세고, 로드 신선도를 확인하고, 타입화된 상태 보고서를 반환합니다. 구조화된 출력을 사용합니다.
- 스키마 드리프트 탐지기(Schema drift detector): 두 스키마를 비교하고 추가, 제거, 유형 변경된 열의 타입화된 목록을 얻으며, 파괴적 변경은 플래그 됩니다. 배포 게이트에 직접 연결됩니다.
- AI 쿼리 최적화기(AI query optimizer): 느린 쿼리를 붙여넣으면 다중 턴 에이전트 세션이 첫 번째 턴에서 병목을 진단하고 두 번째 턴에서 최적화된 재작성을 생성하며, 턴 사이의 전체 컨텍스트를 유지합니다.
각 도구는 자연어 프롬프트를 SDK로 보내고 에이전트가 SQL 작업을 하게 하는 작은 API 라우트입니다. 앱은 쿼리를 하드코딩하지 않습니다.
사전 요구 사항
- Node.js 22 이상.
- Cortex Code CLI 설치:
curl -LsS https://ai.snowflake.com/static/cc-scripts/install.sh | sh - Snowflake CLI 연결 설정을 통해 구성된 Snowflake 연결(일반적으로
~/.snowflake/connections.toml, 기존 설정에~/.snowflake/config.toml도 지원):[my-connection] account = "myorg-myaccount" user = "myuser" authenticator = "externalbrowser" - 코파일럿이 가리킬 스키마. 계정의 아무 스키마나 작동합니다. 선택 사항: demo 데이터 생성은 의도적으로 혼합된 상태의 테이블이 있는 작은 스키마를 만드는 스크립트를 포함합니다.
프로젝트 설정
프로젝트를 만들고 SDK를 설치합니다.
npx create-next-app@latest data-copilot --ts --app --no-src-dir
cd data-copilot
npm install cortex-code-agent-sdk
SDK는 Snowflake CLI 연결을 통해 인증합니다. 명명된 연결을 가리키지 않으면 기본 연결을 사용합니다. 앱은 환경 변수에서 연결 이름을 읽어 같은 코드가 로컬, Cortex Code Desktop, Snowpark Container Services에서 작동하게 합니다. 지원되는 두 변수는 이 순서로 확인됩니다.
SNOWFLAKE_DEFAULT_CONNECTION_NAME: 기본 변수. Cortex Code Desktop과 Snowpark Container Services가 자동으로 설정하므로 이미 설정되어 있을 수 있습니다.SNOWFLAKE_CONNECTION_NAME: 앱을 다른 명명된 연결로 가리키기 위해 설정할 수 있는 재정의.
둘 다 없으면 SDK는 Snowflake CLI 기본 연결로 폴백합니다.
# Optional: only needed to override the automatically detected connection
export SNOWFLAKE_DEFAULT_CONNECTION_NAME=my-connection
웹 UI를 만들기 전에 한 파일 스크립트로 SDK가 Snowflake에 닿는지 확인하세요. scripts/smoke-test.ts로 저장하고 npx tsx scripts/smoke-test.ts로 실행합니다.
// scripts/smoke-test.ts
import { query } from "cortex-code-agent-sdk";
for await (const message of query({
prompt: "List the tables in the SALES.PUBLIC schema with their row counts.",
options: { allowedTools: ["SQL"] }, // auto-approve the SQL tool so the script runs unattended
})) {
if (message.type === "assistant") {
for (const block of message.content) {
if (block.type === "text") process.stdout.write(block.text);
}
}
}
결과가 출력되면 연결된 것입니다. 이 가이드의 나머지는 이 동일한 query() 호출을 API 라우트에 감쌉니다.
프로젝트 구조
완성된 앱은 도구당 하나의 API 라우트, 공유 SDK 헬퍼, 도구당 하나의 UI 페이지를 갖습니다. SDK 코드가 포함된 것은 API 라우트와 헬퍼뿐이고, 페이지는 일반 React입니다.
data-copilot/
├─ app/
│ ├─ api/
│ │ ├─ pipeline-audit/route.ts # structured-output audit, streamed to the UI
│ │ ├─ schema-diff/route.ts # structured-output diff for a deploy gate
│ │ └─ query-optimizer/route.ts # multi-turn diagnose-then-rewrite session
│ ├─ pipeline-audit/page.tsx # UI for each tool (no SDK code)
│ ├─ schema-diff/page.tsx
│ └─ query-optimizer/page.tsx
└─ lib/
└─ sdk.ts # shared SDK options and helpers
선택 사항: demo 데이터 생성
알려진 문제가 있는 스키마에 대해 감사를 재현하려면 워크시트에서 이 스크립트를 실행하세요. 서로 다른 상태(정상, 오래됨, 비어 있음, 높은 NULL)의 테이블 다섯 개를 만듭니다.
CREATE DATABASE IF NOT EXISTS DATA_COPILOT_DEMO;
CREATE SCHEMA IF NOT EXISTS DATA_COPILOT_DEMO.PUBLIC;
USE SCHEMA DATA_COPILOT_DEMO.PUBLIC;
-- Healthy: recently loaded
CREATE OR REPLACE TABLE ORDERS (
order_id INTEGER, customer_id INTEGER, amount DECIMAL(10, 2),
status VARCHAR(20), created_at TIMESTAMP_NTZ, loaded_at TIMESTAMP_NTZ
);
INSERT INTO ORDERS
SELECT seq4() + 1, UNIFORM(1, 5000, RANDOM()),
ROUND(UNIFORM(10, 2000, RANDOM()), 2)::DECIMAL(10, 2),
CASE UNIFORM(1, 3, RANDOM()) WHEN 1 THEN 'PENDING' WHEN 2 THEN 'SHIPPED' ELSE 'DELIVERED' END,
DATEADD('minute', -UNIFORM(1, 1440, RANDOM()), CURRENT_TIMESTAMP()),
DATEADD('minute', -UNIFORM(5, 90, RANDOM()), CURRENT_TIMESTAMP())
FROM TABLE(GENERATOR(ROWCOUNT => 100000));
-- Healthy: fully populated
CREATE OR REPLACE TABLE CUSTOMERS (
customer_id INTEGER, name VARCHAR(100), email VARCHAR(200),
region VARCHAR(50), plan_tier VARCHAR(20), loaded_at TIMESTAMP_NTZ
);
INSERT INTO CUSTOMERS
SELECT seq4() + 1, 'Customer_' || LPAD(seq4() + 1, 5, '0'),
'user' || (seq4() + 1) || '@example.com',
CASE UNIFORM(1, 4, RANDOM()) WHEN 1 THEN 'North' WHEN 2 THEN 'South' WHEN 3 THEN 'East' ELSE 'West' END,
CASE UNIFORM(1, 3, RANDOM()) WHEN 1 THEN 'free' WHEN 2 THEN 'pro' ELSE 'enterprise' END,
DATEADD('minute', -UNIFORM(10, 120, RANDOM()), CURRENT_TIMESTAMP())
FROM TABLE(GENERATOR(ROWCOUNT => 10000));
-- Stale: last loaded about 3 days ago
CREATE OR REPLACE TABLE PRODUCT_EVENTS (
event_id INTEGER, product_id INTEGER, event_type VARCHAR(50),
session_id VARCHAR(50), user_id INTEGER, loaded_at TIMESTAMP_NTZ
);
INSERT INTO PRODUCT_EVENTS
SELECT seq4() + 1, UNIFORM(1, 500, RANDOM()),
CASE UNIFORM(1, 4, RANDOM()) WHEN 1 THEN 'view' WHEN 2 THEN 'click' WHEN 3 THEN 'add_to_cart' ELSE 'purchase' END,
'sess_' || UNIFORM(10000, 99999, RANDOM()), UNIFORM(1, 10000, RANDOM()),
DATEADD('hour', -UNIFORM(80, 96, RANDOM()), CURRENT_TIMESTAMP())
FROM TABLE(GENERATOR(ROWCOUNT => 25000));
-- Empty: pipeline never populated it
CREATE OR REPLACE TABLE SESSION_METRICS (
session_id VARCHAR(100), user_id INTEGER, duration_seconds INTEGER,
page_views INTEGER, bounce BOOLEAN, loaded_at TIMESTAMP_NTZ
);
-- High nulls: about a third of campaign_id and channel are NULL
CREATE OR REPLACE TABLE AD_SPEND (
spend_id INTEGER, campaign_id VARCHAR(50), channel VARCHAR(50),
impressions INTEGER, clicks INTEGER, cost DECIMAL(10, 2), loaded_at TIMESTAMP_NTZ
);
INSERT INTO AD_SPEND
SELECT seq4() + 1,
CASE WHEN UNIFORM(1, 3, RANDOM()) = 1 THEN NULL ELSE 'campaign_' || LPAD(UNIFORM(1, 30, RANDOM()), 3, '0') END,
CASE WHEN UNIFORM(1, 3, RANDOM()) = 1 THEN NULL
ELSE CASE UNIFORM(1, 4, RANDOM()) WHEN 1 THEN 'email' WHEN 2 THEN 'social' WHEN 3 THEN 'search' ELSE 'display' END END,
UNIFORM(100, 50000, RANDOM()), UNIFORM(5, 2000, RANDOM()),
ROUND(UNIFORM(50, 5000, RANDOM()), 2)::DECIMAL(10, 2),
DATEADD('hour', -UNIFORM(1, 6, RANDOM()), CURRENT_TIMESTAMP())
FROM TABLE(GENERATOR(ROWCOUNT => 8000));
핵심 SDK 사용법 살펴보기
공유 옵션 헬퍼
모든 도구가 SDK를 같은 방식으로 구성하므로 lib/sdk.ts에 중앙화할 가치가 있습니다. 앱은 에이전트를 무인으로 실행합니다(웹 요청에서 각 도구 호출을 승인할 사람이 없음). 그래서 permissionMode: "bypassPermissions"와 필수 allowDangerouslySkipPermissions 안전 플래그로 도구를 자동 승인하고, 외부 MCP(Model Context Protocol) 서버를 비활성화하며, 환경에서 연결 이름을 가져옵니다.
같은 파일에는 에이전트의 스트리밍 메시지에서 텍스트와 SQL을 빼내는 두 개의 작은 헬퍼도 있습니다. 도구들이 브라우저에 에이전트 작업을 표시하는 데 재사용합니다.
// lib/sdk.ts
import type {
CortexCodeSessionOptions,
ContentBlock,
TextBlock,
ToolUseBlock,
} from "cortex-code-agent-sdk";
export function sdkOptions(
override?: Partial<CortexCodeSessionOptions>,
): CortexCodeSessionOptions {
// Cortex Code Desktop and Snowpark Container Services set SNOWFLAKE_DEFAULT_CONNECTION_NAME.
const connection =
process.env.SNOWFLAKE_DEFAULT_CONNECTION_NAME ||
process.env.SNOWFLAKE_CONNECTION_NAME;
return {
permissionMode: "bypassPermissions",
allowDangerouslySkipPermissions: true,
noMcp: true,
...(connection ? { connection } : {}),
...override,
};
}
// Plain text the agent streams back (its reasoning and explanations).
export function textBlocks(content: ContentBlock[]): TextBlock[] {
return content.filter((b): b is TextBlock => b.type === "text");
}
// Each SQL statement the agent runs, so the UI can show its work.
export function sqlToolUses(content: ContentBlock[]): string[] {
return content
.filter((b): b is ToolUseBlock => b.type === "tool_use" && b.name === "SQL")
.map((b) => {
const input = b.input as Record<string, unknown>;
// The SQL tool input uses "command" or "query" depending on CLI version.
return (input.command ?? input.query ?? JSON.stringify(input)) as string;
});
}
경고:
permissionMode: "bypassPermissions"은 모든 도구를 프롬프트 없이 실행합니다. 자신의 서버나 CI/CD 작업 같은 신뢰되는 샌드박스 환경에서만 사용하세요. 인간 감독이 필요한 워크플로에는allowedTools,disallowedTools또는canUseTool콜백을 대신 사용하세요. 승인 및 사용자 입력 처리를 참고하세요.
도구 1: 파이프라인 상태 감사
감사는 하나의 프롬프트와 JSON Schema를 보내고 에이전트가 SQL을 작성·실행하게 합니다. 요청이 outputFormat을 전달하므로 최종 결과가 스키마에 대해 검증되어, 라우트가 파싱해야 할 자유 텍스트 대신 타입화된 보고서를 받습니다. 에이전트가 작업하는 동안 라우트는 lib/sdk.ts의 헬퍼를 사용해 그 텍스트와 각 SQL 문을 브라우저로 스트리밍합니다.
// app/api/pipeline-audit/route.ts (SDK essentials)
import { query } from "cortex-code-agent-sdk";
import { sdkOptions, textBlocks, sqlToolUses } from "@/lib/sdk";
const AUDIT_OUTPUT_SCHEMA = {
type: "object",
properties: {
tables: {
type: "array",
items: {
type: "object",
properties: {
name: { type: "string" },
rowCount: { type: "number" },
lastLoaded: { type: ["string", "null"] },
hoursSinceLoad: { type: ["number", "null"] },
status: { type: "string", enum: ["healthy", "stale", "empty", "high_nulls"] },
recommendation: { type: "string" },
},
required: ["name", "rowCount", "lastLoaded", "hoursSinceLoad", "status", "recommendation"],
},
},
total: { type: "number" },
healthy: { type: "number" },
issues: { type: "number" },
},
required: ["tables", "total", "healthy", "issues"],
};
const prompt = `Audit schema ${database}.${schema} in Snowflake. For each BASE TABLE:
1. Get the row count from INFORMATION_SCHEMA.TABLES (ROW_COUNT column)
2. Find the most recent timestamp column: LOADED_AT, UPDATED_AT, CREATED_AT, or any TIMESTAMP/DATE column in INFORMATION_SCHEMA.COLUMNS
3. For tables with a timestamp column, query MAX of that column to compute hours since last load
4. Classify each table: "empty" (0 rows), "stale" (last load > 24h ago), "high_nulls" (if you detect high null rates), or "healthy"
5. Write a short, specific recommendation for each non-healthy table
Return the full health report as structured JSON.`;
for await (const event of query({
prompt,
options: sdkOptions({
allowedTools: ["SQL"],
// The SDK types this field as Record<string, unknown>, so cast the const schema.
outputFormat: { type: "json_schema", schema: AUDIT_OUTPUT_SCHEMA as Record<string, unknown> },
}),
})) {
if (event.type === "system" && event.subtype === "init") {
// The init event reports which model the agent is using.
console.log(`Agent started, model: ${event.model}`);
}
if (event.type === "assistant") {
for (const block of textBlocks(event.content)) {
// stream block.text to the UI
}
for (const sql of sqlToolUses(event.content)) {
// stream each SQL statement the agent runs to the UI
}
}
if (event.type === "result" && event.subtype === "success") {
// structured_output is typed unknown, so cast it to the shape you asked for.
const report = event.structured_output as {
tables: unknown[];
total: number;
healthy: number;
issues: number;
};
console.log(`${report.issues} issue(s) across ${report.total} tables`);
}
}
from pydantic import BaseModel
from cortex_code_agent_sdk import query, AssistantMessage, ResultMessage, CortexCodeAgentOptions
class TableHealth(BaseModel):
name: str
row_count: int
last_loaded: str | None
hours_since_load: float | None
status: str # "healthy" | "stale" | "empty" | "high_nulls"
recommendation: str
class PipelineReport(BaseModel):
tables: list[TableHealth]
total: int
healthy: int
issues: int
async for msg in query(
prompt=f"Audit schema {database}.{schema}: flag stale (>24h), empty, and high-null tables",
options=CortexCodeAgentOptions(
connection="my-connection",
allowed_tools=["SQL"],
output_format={"type": "json_schema", "schema": PipelineReport.model_json_schema()},
),
):
if isinstance(msg, ResultMessage) and msg.structured_output:
report = PipelineReport(**msg.structured_output)
print(f"{report.issues} issue(s) across {report.total} tables")
도구 2: 스키마 드리프트 탐지기
드리프트 탐지기도 같은 구조화된 출력 패턴을 사용하지만, 타입화된 결과가 결정을 구동합니다: 에이전트가 파괴적 변경을 찾으면 앱이 배포를 실패시킬 수 있습니다. 이로써 자연어 요청이 결정적 게이트로 바뀝니다.
// app/api/schema-diff/route.ts (SDK essentials)
import { query } from "cortex-code-agent-sdk";
import { sdkOptions } from "@/lib/sdk";
const SCHEMA_DIFF_OUTPUT = {
type: "object",
properties: {
sourceSchema: { type: "string" },
targetSchema: { type: "string" },
tablesCompared: { type: "number" },
addedTables: { type: "array", items: { type: "string" } },
removedTables: { type: "array", items: { type: "string" } },
changes: {
type: "array",
items: {
type: "object",
properties: {
tableName: { type: "string" },
columnName: { type: "string" },
changeType: { type: "string", enum: ["added", "removed", "type_changed"] },
sourceType: { type: ["string", "null"] },
targetType: { type: ["string", "null"] },
isBreaking: { type: "boolean" },
},
required: ["tableName", "columnName", "changeType", "sourceType", "targetType", "isBreaking"],
},
},
breakingChanges: { type: "number" },
},
required: ["sourceSchema", "targetSchema", "tablesCompared", "addedTables", "removedTables", "changes", "breakingChanges"],
};
const prompt = `Compare the column schemas of every table in ${srcDb}.${srcSchema} with its counterpart in ${tgtDb}.${tgtSchema}.
Use INFORMATION_SCHEMA.COLUMNS to get column names and data types for both schemas.
Identify columns that were added, removed, or had their type changed between source and target.
Flag type changes as breaking if they narrow the type (for example, a wider VARCHAR to a narrower VARCHAR, or NUMBER with reduced precision).
Also identify tables that exist in one schema but not the other.
Return a complete structured diff.`;
for await (const event of query({
prompt,
options: sdkOptions({
allowedTools: ["SQL"],
// The SDK types this field as Record<string, unknown>, so cast the const schema.
outputFormat: { type: "json_schema", schema: SCHEMA_DIFF_OUTPUT as Record<string, unknown> },
}),
})) {
if (event.type === "result" && event.subtype === "success") {
// structured_output is typed unknown, so cast it to the shape you asked for.
const diff = event.structured_output as { breakingChanges: number };
if (diff.breakingChanges > 0) {
process.exit(1); // In CI, block the deploy
}
}
}
from pydantic import BaseModel
from cortex_code_agent_sdk import query, ResultMessage, CortexCodeAgentOptions
class ColumnChange(BaseModel):
table_name: str
column_name: str
change_type: str # "added" | "removed" | "type_changed"
source_type: str | None
target_type: str | None
is_breaking: bool
class SchemaDiff(BaseModel):
source_schema: str
target_schema: str
tables_compared: int
added_tables: list[str]
removed_tables: list[str]
changes: list[ColumnChange]
breaking_changes: int
async for msg in query(
prompt=(
"Compare schemas STAGING.PUBLIC and PROD.PUBLIC using INFORMATION_SCHEMA. "
"Flag removed or type-narrowed columns as breaking."
),
options=CortexCodeAgentOptions(
connection="my-connection",
allowed_tools=["SQL"],
output_format={"type": "json_schema", "schema": SchemaDiff.model_json_schema()},
),
):
if isinstance(msg, ResultMessage) and msg.structured_output:
diff = SchemaDiff(**msg.structured_output)
if diff.breaking_changes > 0:
raise SystemExit(f"Deploy blocked: {diff.breaking_changes} breaking change(s)")
도구 3: AI 쿼리 최적화기
최적화기는 컨텍스트를 공유하는 두 턴이 필요합니다: 먼저 진단, 다음 재작성. query() 대신 세션을 사용해 두 번째 프롬프트가 첫 번째에서 에이전트가 배운 모든 것에 의존하게 합니다. 또한 시스템 프롬프트를 추가해 에이전트를 Snowflake 특화 조언으로 유도합니다.
// app/api/query-optimizer/route.ts (SDK essentials)
import { createCortexCodeSession } from "cortex-code-agent-sdk";
import { sdkOptions, textBlocks, sqlToolUses } from "@/lib/sdk";
const session = await createCortexCodeSession(
sdkOptions({
allowedTools: ["SQL"],
appendSystemPrompt:
"You are a Snowflake SQL expert. When diagnosing queries, be specific about " +
"micro-partition pruning, clustering keys, partition pruning ratios, and warehouse " +
"sizing. Reference Snowflake-native features by name.",
}),
);
// A session reports its model through initializationResult() instead of an init event.
const init = await session.initializationResult();
const model = (init as Record<string, unknown>).model as string;
console.log(`Agent started, model: ${model}`);
// Stream one turn's text and SQL to the UI. Both turns are handled the same way,
// so the agent's rewrite in turn 2 is streamed just like the diagnosis in turn 1.
async function streamTurn() {
for await (const event of session.stream()) {
if (event.type === "assistant") {
for (const block of textBlocks(event.content)) { /* stream the text to the UI */ }
for (const ranSql of sqlToolUses(event.content)) { /* stream the SQL the agent ran */ }
}
if (event.type === "result") break;
}
}
// Turn 1: diagnose the bottleneck
await session.send(
`Diagnose this Snowflake SQL query. Identify the main performance bottleneck and explain why it's slow:\n\n\`\`\`sql\n${sql}\n\`\`\``,
);
await streamTurn();
// Turn 2: rewrite. The agent still has the full diagnosis in context.
await session.send(
"Now produce the full optimized rewrite of that query with inline comments explaining each key change.",
);
await streamTurn();
await session.close();
from cortex_code_agent_sdk import CortexCodeSDKClient, CortexCodeAgentOptions, AssistantMessage
system_prompt = {"type": "preset", "append": (
"You are a Snowflake SQL expert. When diagnosing queries, be specific about "
"micro-partition pruning, clustering keys, partition pruning ratios, and warehouse "
"sizing. Reference Snowflake-native features by name."
)}
async with CortexCodeSDKClient(
CortexCodeAgentOptions(connection="my-connection", allowed_tools=["SQL"], system_prompt=system_prompt)
) as client:
# Turn 1: diagnose
await client.query(f"Diagnose this Snowflake SQL query. Identify the main performance bottleneck and explain why it's slow:\n\n```sql\n{slow_query}\n```")
async for msg in client.receive_response():
if isinstance(msg, AssistantMessage):
for block in msg.content:
if hasattr(block, "text"):
print(block.text, end="")
# Turn 2: rewrite. The agent still has full context from turn 1.
await client.query("Now produce the full optimized rewrite of that query with inline comments explaining each key change.")
async for msg in client.receive_response():
if isinstance(msg, AssistantMessage):
for block in msg.content:
if hasattr(block, "text"):
print(block.text, end="")
실행하기
npm run dev
http://localhost:3000 을 열고 도구를 선택한 뒤 원하는 것을 설명하세요. 이 섹션의 나머지는 각 도구가 만드는 출력 종류를 보여줍니다.
파이프라인 상태 감사 출력
감사가 실행되는 동안 라우트가 에이전트의 추론과 작성한 각 SQL 문을 스트리밍하므로 브라우저가 작업 중인 것을 보여줍니다.
Agent started, model: auto
Auditing DATA_COPILOT_DEMO.PUBLIC: fetching tables, row counts, and timestamp columns
SQL: SELECT table_name, row_count FROM DATA_COPILOT_DEMO.INFORMATION_SCHEMA.TABLES WHERE table_schema = 'PUBLIC'
SQL: SELECT MAX(loaded_at) FROM DATA_COPILOT_DEMO.PUBLIC.PRODUCT_EVENTS
최종 결과 이벤트는 타입화된 PipelineReport를 전달합니다. 앱은 요약 개수와 테이블당 한 행으로 렌더링합니다.
| Table | Rows | Last loaded | Status | Recommendation |
|---|---|---|---|---|
| ORDERS | 100,000 | ~1 hour ago | healthy | None |
| CUSTOMERS | 10,000 | ~1 hour ago | healthy | None |
| PRODUCT_EVENTS | 25,000 | ~3 days ago | stale | Last load is about 3 days old; run the ingestion job to pick up new events. |
| SESSION_METRICS | 0 | Never | empty | Table has no rows; confirm the pipeline populates it or drop it. |
| AD_SPEND | 8,000 | ~2 hours ago | high_nulls | campaign_id and channel are about 33% NULL; check the source mapping. |
두 테이블은 정상이고 세 개는 주의가 필요합니다: PRODUCT_EVENTS는 오래됐고, SESSION_METRICS는 비어 있으며, AD_SPEND는 높은 NULL 비율입니다.
스키마 드리프트 탐지기 출력
스테이징 스키마와 프로덕션을 비교하면 배포 게이트가 작동할 수 있는 타입화된 SchemaDiff를 반환합니다.
{
"sourceSchema": "STAGING.PUBLIC",
"targetSchema": "PROD.PUBLIC",
"tablesCompared": 12,
"addedTables": ["FEATURE_FLAGS"],
"removedTables": [],
"changes": [
{ "tableName": "ORDERS", "columnName": "discount_pct", "changeType": "added", "sourceType": null, "targetType": "NUMBER(5,2)", "isBreaking": false },
{ "tableName": "CUSTOMERS", "columnName": "email", "changeType": "type_changed", "sourceType": "VARCHAR(255)", "targetType": "VARCHAR(100)", "isBreaking": true },
{ "tableName": "SESSIONS", "columnName": "user_agent", "changeType": "removed", "sourceType": "VARCHAR(1000)", "targetType": null, "isBreaking": true }
],
"breakingChanges": 2
}
breakingChanges가 0보다 크므로 이전 섹션의 CI 단계가 0이 아닌 값으로 종료되어 배포를 차단합니다.
AI 쿼리 최적화기 출력
느린 쿼리를 붙여넣으세요.
SELECT
u.userid,
u.firstname || ' ' || u.lastname AS buyer_name,
u.city,
COUNT(s.salesid) AS total_purchases,
SUM(s.qtysold * s.pricepaid) AS total_spent,
AVG(s.pricepaid) AS avg_ticket_price
FROM SAMPLES.TICKIT.SALES s
JOIN SAMPLES.TICKIT.USERS u ON s.buyerid = u.userid
JOIN SAMPLES.TICKIT.EVENT e ON s.eventid = e.eventid
WHERE s.saletime >= DATEADD('month', -6, CURRENT_TIMESTAMP())
AND e.catid IN (1, 2, 3)
GROUP BY 1, 2, 3
ORDER BY total_spent DESC
LIMIT 100;
턴 1이 병목을 진단합니다. 에이전트가 테이블 메타데이터를 검사하고 SALES에 대한 전체 마이크로 파티션 스캔이라는 기본 문제를 보고합니다.
SALES: 172,456 rows | cluster_by: (none) | automatic_clustering: OFF
s.saletime >= DATEADD('month', -6, CURRENT_TIMESTAMP()) 조건자는 SALES에 클러스터링 키가 없어 어떤 마이크로 파티션도 잘라낼 수 없으므로 Snowflake가 모든 파티션을 스캔합니다. 에이전트는 또한 부차적 병목을 플래그합니다: e.catid IN (1, 2, 3) 필터는 EVENT에 있고 조인이 완료된 후에만 SALES 행을 제거합니다.
턴 2는 턴 1의 진단을 재사용해 재작성을 생성합니다. EVENT를 CTE에서 미리 필터링하고, 가장 작은 필터링된 집합이 쿼리를 구동하도록 조인을 재정렬하며, 위치 별칭 대신 원시 열로 그룹화합니다.
-- OPTIMIZED REWRITE: SAMPLES.TICKIT buyer summary
WITH eligible_events AS (
-- [1] Pre-filter EVENT before any join to SALES, so only matching eventids reach the SALES join.
SELECT eventid
FROM SAMPLES.TICKIT.EVENT
WHERE catid IN (1, 2, 3)
)
SELECT
u.userid,
-- [4] Concatenation deferred until after grouping (runs on ~100 result rows, not every input row).
u.firstname || ' ' || u.lastname AS buyer_name,
u.city,
COUNT(s.salesid) AS total_purchases,
SUM(s.qtysold * s.pricepaid) AS total_spent,
AVG(s.pricepaid) AS avg_ticket_price
FROM SAMPLES.TICKIT.SALES s
JOIN eligible_events e ON s.eventid = e.eventid -- [2] join the smallest filtered table first
JOIN SAMPLES.TICKIT.USERS u ON s.buyerid = u.userid -- [3] USERS joins already-filtered SALES rows
WHERE s.saletime >= DATEADD('month', -6, CURRENT_TIMESTAMP()) -- [5] benefits from CLUSTER BY (saletime)
GROUP BY u.userid, u.firstname, u.lastname, u.city
ORDER BY total_spent DESC
LIMIT 100;
재작성은 에이전트가 전제 조건으로 지적하는 일회성 클러스터링 변경 후에만 완전한 이점을 제공합니다.
-- Snowflake reclusters asynchronously in the background.
ALTER TABLE SAMPLES.TICKIT.SALES CLUSTER BY (saletime);
확장하기
같은 구성 요소가 더 많은 자동화를 지원합니다.
- CI/CD에서 헤드리스 실행: 인간 개입 없이 예약 작업에서 감사 또는 드리프트 확인을 구동. 폭주 작업을
maxTurns로 제한. - 에이전트가 실행한 모든 쿼리 기록: 각 SQL 문을 규정 준수 감사 로그에 추가하는
PostToolUse훅 추가. Hooks 참고. - 로컬 파일과 Snowflake 결합: 에이전트에
SQL과 함께Read,Glob,Write도구를 주고cwd를 dbt 프로젝트로 가리켜 모델을 실제 스키마와 교차 참조하게 함.
다음 단계
- 구조화된 출력: 에이전트 워크플로에서 검증된 JSON 반환
- 다중 턴 세션 및 스트리밍 입력: 여러 교환에 걸쳐 컨텍스트 유지
- 승인 및 사용자 입력 처리: 에이전트가 사용할 수 있는 도구 제어
- Hooks: 에이전트 수명 주기의 주요 지점에서 사용자 정의 코드 실행
- TypeScript SDK 참조: 전체 TypeScript API 참조
- Python SDK 참조: 전체 Python API 참조