import os
from contextlib import contextmanager # contextmanager 임포트
from typing import Generator
from dotenv import load_dotenv
from sqlalchemy.ext.asyncio import create_async_engine, async_sessionmaker, AsyncSession
from sqlalchemy.orm import sessionmaker, declarative_base, Session # Session 타입 힌트 추가
load_dotenv()
# 환경 변수 로딩 시 변수명 축약 대신 명시적인 이름 사용 권장
MYSQL_USER = os.environ.get("MYSQL_USER")
MYSQL_PASSWORD = os.environ.get("MYSQL_PASSWORD")
MYSQL_HOST = os.environ.get("MYSQL_ALEMBIC_HOST") # alembic 전용 호스트가 아니라면 MYSQL_HOST로 통일 가능
MYSQL_PORT = os.environ.get("MYSQL_ALEMBIC_PORT")
MYSQL_DB = os.environ.get("MYSQL_DATABASE")
# MySQL 연결 설정 (사용자, 비밀번호, 호스트, DB이름 수정)
# 예: mysql+aiomysql://user:password@localhost:3306/dbname
SQLALCHEMY_DATABASE_URL = f"mysql+aiomysql://{MYSQL_USER}:{MYSQL_PASSWORD}@{MYSQL_HOST}:{MYSQL_PORT}/{MYSQL_DB}"
# SQLite 사용 시: SQLALCHEMY_DATABASE_URL = "sqlite:///./myapi.db"
# create_engine 설정: MySQL 연결 시 pool_pre_ping은 좋은 관례입니다.
engine = create_async_engine(
SQLALCHEMY_DATABASE_URL,
echo=True,
pool_pre_ping=True, # 연결 유효성 체크 (MySQL 연결 끊김 방지)
pool_recycle=3600, # 연결 재사용 시간 설정
# SQLite 사용 시: connect_args={"check_same_thread": False}
)
# 세션 생성을 위한 팩토리
async_session = async_sessionmaker(
bind=engine,
class_=AsyncSession,
expire_on_commit=False
)
# Base 모델 클래스 생성
Base = declarative_base()
# 참고: Alembic 사용 시 필요한 주석들은 제거하거나 별도의 alembic config 파일로 옮기는 것이 깔끔합니다.
BOOTCAMP
본 공간은 인공지능(AI)이 도출한 방대한 지식을 일목요연하게 큐레이션하여 기록하는 블로그입니다. 기술적 한계로 인해 모든 정보의 완벽한 정확성을 보장하기는 어려우므로, 최종적인 판단과 책임은 독자 본인에게 있음을 정중히 안내해 드립니다.
2026년 7월 21일 화요일
database.py
Cargo.toml
# curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh
# rustup update
# cargo update
[package]
name = "sheets"
version = "0.1.0"
edition = "2024"
default-run = "sheets"
[[bin]]
name = "sheets"
path = "src/main.rs"
[profile.dev.package.backtrace]
opt-level = 3
[dependencies]
tokio = { version = "1", features = ["full"] }
futures = "0.3"
itertools = "0.12"
serde = "1"
serde_json = "1"
rustls = { version = "0.23", default-features = false, features = ["aws-lc-rs", "logging", "tls12"] }
google-sheets4 = "*"
yup-oauth2 = "12"
Inflector = "0.11.4"
src/main.rs
use std::path::Path;
use tokio::fs::File;
use tokio::io::AsyncWriteExt; // 비동기 쓰기를 위한 트레이트
mod gspread;
use crate::gspread::{SpreadsheetExporter, create_hub};
/// 생성된 SQL 문을 비동기적으로 .sql 파일에 저장
async fn save_sql_to_file_async(sql: &str, file_name: &str) -> tokio::io::Result<()> {
// 1. 파일 경로 설정
let path = Path::new(file_name);
// 2. 비동기 파일 생성 (Tokio fs 사용)
let mut file = File::create(path).await?;
// 3. 비동기 데이터 쓰기 (전체 데이터를 안전하게 기록)
file.write_all(sql.as_bytes()).await?;
// 4. 버퍼 비우기 및 동기화 (권장 관례)
file.flush().await?;
println!("✅ [Async] SQL 파일 저장 완료: {:?}", path.display());
Ok(())
}
#[tokio::main]
async fn main() -> tokio::io::Result<()> {
rustls::crypto::aws_lc_rs::default_provider()
.install_default()
.expect("Failed to install crypto provider");
let hub = create_hub("./sheets-xxxxxx-xxxx.json").await;
let exporter =
SpreadsheetExporter::new(hub, "xxxxxx-xx-xxxxxx".into());
let schema = exporter.export_schema().await;
save_sql_to_file_async(&schema, "output_schema.sql").await?;
let data = exporter.export_data().await;
save_sql_to_file_async(&data, "output_data.sql").await?;
println!("{}\n{}", schema, data);
Ok(())
}
src/gspread.rs
//extern crate hyper;
//extern crate hyper_rustls;
//extern crate hyper_util;
extern crate google_sheets4 as sheets4;
use sheets4::{Sheets, hyper_rustls, hyper_util, yup_oauth2};
use serde_json::Value;
use std::collections::HashMap;
use futures::future::join_all;
use hyper_util::client::legacy::connect::HttpConnector;
use hyper_rustls::HttpsConnector;
use inflector::Inflector; // 패키지: id = "0.11" (네이밍 변환용)
pub type SheetsHub = Sheets>;
//#[derive(Default, Debug, Clone, PartialEq)]
pub struct SpreadsheetExporter {
hub: SheetsHub,
spreadsheet_id: String,
}
impl SpreadsheetExporter {
pub fn new(hub: SheetsHub, spreadsheet_id: String) -> Self {
Self { hub, spreadsheet_id }
}
// Value 타입을 String으로 안전하게 변환하는 헬퍼 함수
fn val_to_string(val: &Value) -> String {
match val {
Value::String(s) => s.clone(),
Value::Number(n) => n.to_string(),
Value::Bool(b) => b.to_string(),
_ => val.to_string().replace('"', ""),
}
}
// 시트에서 가져온 Vec>를 정제된 HashMap 리스트로 변환
fn clean_rows(&self, values: &Vec>) -> Vec> {
if values.is_empty() { return vec![]; }
let headers: Vec = values[0].iter().map(Self::val_to_string).collect();
let mut result = Vec::new();
for row in values.iter().skip(1) {
if row.is_empty() || Self::val_to_string(&row[0]).trim().is_empty() { break; }
let mut row_map = HashMap::new();
for (i, header) in headers.iter().enumerate() {
let val = row.get(i).map(Self::val_to_string).unwrap_or_default();
row_map.insert(header.clone(), val);
}
result.push(row_map);
}
result
}
fn split_column(&self, raw_text: &str) -> String {
raw_text.split(',').map(|s| format!("`{}`", s.trim())).collect::>().join(", ")
}
fn format_sql_value(&self, val: &str) -> String {
let t = val.trim();
if t.is_empty() || t.to_uppercase() == "NULL" { return "NULL".to_string(); }
if t.parse::().is_ok() { t.to_string() } else { format!("'{}'", t.replace("'", "''")) }
}
async fn export_schema_unit(&self, row: HashMap) -> String {
let schema = row.get("SCHEMA").cloned().unwrap_or_default();
let table_name = row.get("TABLE_NAME").cloned().unwrap_or_default();
let engine = row.get("ENGINE").cloned().unwrap_or_default();
let auto_increment = row.get("AUTO_INCREMENT").cloned().unwrap_or_default();
let charset = row.get("CHARSET").cloned().unwrap_or_default();
let collate = row.get("COLLATE").cloned().unwrap_or_default();
let ss_id = match row.get("SPREADSHEET_ID") {
Some(id) => id,
None => return format!("-- Error: Missing ID for {}", table_name),
};
let result = self.hub.spreadsheets().values_batch_get(ss_id)
.add_ranges("COLUMN!A:Z")
.add_ranges("INDEX!A:Z")
.add_ranges("FOREIGN_KEY!A:Z")
.doit().await;
let (_, response) = match result {
Ok(res) => res,
Err(e) => return format!("-- API Error in {}: {}\n", table_name, e),
};
let v_ranges = response.value_ranges.unwrap_or_default();
let mut all_defs = Vec::new();
// 1. COLUMN 처리
if let Some(vals) = v_ranges.get(0).and_then(|r| r.values.as_ref()) {
for r in self.clean_rows(vals) {
let f = r.get("FIELD").cloned().unwrap_or_default();
let t = r.get("TYPE").cloned().unwrap_or_default();
let s = r.get("SIGNED").cloned().unwrap_or_default();
let l = r.get("LENGTH").filter(|s| !s.is_empty()).map_or("".to_string(), |v| format!("({})", v));
let n = if ["NOT", "NO", "FALSE", "0"].contains(&r.get("NULL").unwrap_or(&"".into()).to_uppercase().as_str()) { "NOT NULL" } else { "" };
let d = r.get("DEFAULT").filter(|s| !s.is_empty()).map_or("".to_string(), |v| format!("DEFAULT '{}'", v.replace("'", "''")));
let c = r.get("COMMENT").filter(|s| !s.is_empty()).map_or("".to_string(), |v| format!("COMMENT '{}'", v.replace("'", "''")));
all_defs.push(format!(" `{}` {} {}{} {} {} {} {}", f.to_snake_case(), t, s, l, n, d, r.get("EXTRA").unwrap_or(&"".into()), c).trim().to_string());
}
}
// 2. INDEX 처리
if let Some(vals) = v_ranges.get(1).and_then(|r| r.values.as_ref()) {
for r in self.clean_rows(vals) {
let itype = r.get("TYPE").unwrap_or(&"".into()).to_uppercase();
let cols = self.split_column(r.get("COLUMN").unwrap_or(&"".into()));
if itype == "PRIMARY" { all_defs.push(format!("PRIMARY KEY ({})", cols.to_snake_case())); }
else {
let prefix = if itype == "KEY" { "".into() } else { format!("{} ", itype) };
all_defs.push(format!("{}KEY `{}` ({})", prefix, r.get("NAME").unwrap_or(&"".into()).to_snake_case(), cols));
}
}
}
// 3. FOREIGN KEY 처리
if let Some(vals) = v_ranges.get(2).and_then(|r| r.values.as_ref()) {
for r in self.clean_rows(vals) {
let mut fk = format!("CONSTRAINT `{}` FOREIGN KEY (`{}`) REFERENCES `{}` (`{}`)",
r.get("NAME").unwrap_or(&"".to_string()).to_snake_case(),
r.get("COLUMN").unwrap_or(&"".to_string()).to_snake_case(),
r.get("REF_TABLE").unwrap_or(&"".to_string()).to_plural().to_snake_case(),
r.get("REF_COLUMN").unwrap_or(&"".to_string()).to_snake_case());
if let Some(u) = r.get("ON_UPDATE").filter(|s| !s.is_empty()) { fk.push_str(&format!(" ON UPDATE {}", u)); }
if let Some(d) = r.get("ON_DELETE").filter(|s| !s.is_empty()) { fk.push_str(&format!(" ON DELETE {}", d)); }
all_defs.push(fk);
}
}
let schema_name = schema.to_snake_case();
let table_name = table_name.to_plural().to_snake_case();
format!("CREATE DATABASE IF NOT EXISTS {};\nUSE {};\nDROP TABLE IF EXISTS `{}`;\nCREATE TABLE `{}` (\n{}\n) ENGINE={} AUTO_INCREMENT={} DEFAULT CHARSET={} COLLATE={};\n",
schema_name, schema_name, table_name, table_name, all_defs.join(",\n"), engine, auto_increment, charset, collate)
}
async fn export_data_unit(&self, row: HashMap) -> String {
let table_name = row.get("TABLE_NAME").cloned().unwrap_or_default();
let ss_id = row.get("SPREADSHEET_ID").expect("No ID");
let result = self.hub.spreadsheets().values_get(ss_id, "DATA!A:Z").doit().await;
let (_, response) = match result {
Ok(res) => res,
Err(_) => return "".into(),
};
let values = response.values.unwrap_or_default();
if values.len() <= 1 { return "".into(); }
let headers: Vec = values[0].iter().map(Self::val_to_string).collect();
let mut insert_rows = Vec::new();
for r in self.clean_rows(&values) {
let row_vals: Vec = headers.iter()
.map(|h| self.format_sql_value(r.get(h).unwrap_or(&"".into())))
.collect();
insert_rows.push(format!("({})", row_vals.join(", ")));
}
let h_sql = headers.iter().map(|h| format!("`{}`", h.to_snake_case())).collect::>().join(", ");
format!("INSERT INTO `{}` ({})\nVALUES\n{};\n", table_name.to_plural(), h_sql, insert_rows.join(",\n"))
}
async fn execute_parallel(&self, unit_func: F) -> String
where
F: Fn(HashMap) -> Fut,
Fut: std::future::Future
main.py
from database import async_session
from models.answer import Answer
from models import Question
from sqlalchemy import select, insert
from sqlalchemy.exc import SQLAlchemyError
import asyncio
import SpreadsheetExporter
async def create_user_with_transaction(username: str):
# 1. 세션 생성
async with async_session() as session:
# 2. 트랜잭션 시작 (비동기 컨텍스트 매니저)
async with session.begin():
try:
# 데이터 생성 작업
new_user = Answer(username=username)
session.add(new_user)
# 추가 작업 (예: 로그 기록 등)
# 이 블록 안에서 에러가 발생하면 전체 작업이 롤백됩니다.
except Exception as e:
# session.begin()을 사용하면 에러 발생 시 자동 롤백되지만,
# 추가적인 로깅이 필요하면 여기서 처리합니다.
print(f"에러 발생: {e}")
raise
# 블록을 나가면 자동으로 commit 완료
print("트랜잭션이 성공적으로 커밋되었습니다.")
async def manual_transaction_example(username: str):
session = async_session()
try:
# 작업 수행
new_user = Answer(username=username)
session.add(new_user)
# 명시적 커밋
await session.commit()
print("커밋 성공")
except SQLAlchemyError as e:
# 에러 발생 시 명시적 롤백
await session.rollback()
print(f"롤백 실행: {e}")
finally:
# 세션 닫기
await session.close()
async def create_user_and_post(username: str, title: str):
async with async_session() as session:
async with session.begin():
# 1. 사용자 생성
user = Answer(username=username)
session.add(user)
await session.flush() # DB에 임시 반영하여 user.id를 확보
# 2. 해당 사용자의 포스트 생성
post = Question(title=title, owner_id=user.id)
session.add(post)
# 여기서 에러 발생 시 User 생성도 취소됨
async def fast_bulk_insert(user_data_list: list):
async with async_session() as session:
async with session.begin():
# Core의 insert 구문 사용
stmt = insert(Answer).values(user_data_list)
await session.execute(stmt)
import gspread_asyncio
# from google-auth package
from google.oauth2.service_account import Credentials
# First, set up a callback function that fetches our credentials off the disk.
# gspread_asyncio needs this to re-authenticate when credentials expire.
def get_creds():
# To obtain a service account JSON file, follow these steps:
# https://gspread.readthedocs.io/en/latest/oauth2.html#for-bots-using-service-account
#path = os.environ.get('GOOGLE_APPLICATION_CREDENTIALS')
path = "./sheets-xxxxxx-xxxxxxxxxxxx.json"
creds = Credentials.from_service_account_file(path)
scoped = creds.with_scopes([
"https://spreadsheets.google.com/feeds",
"https://www.googleapis.com/auth/spreadsheets",
"https://www.googleapis.com/auth/drive",
])
return scoped
if __name__ == "__main__":
# Execute when the module is not initialized from an import statement.
#main()
# --- 사용 예시 ---
async def main():
# Create an AsyncioGspreadClientManager object which
# will give us access to the Spreadsheet API.
agcm = gspread_asyncio.AsyncioGspreadClientManager(get_creds)
# Always authorize first.
# If you have a long-running program call authorize() repeatedly.
agc = await agcm.authorize()
spreadsheet_id = "xxxxxxxxx-xxx-xxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"
exporter = SpreadsheetExporter.SpreadsheetExporter(agc, spreadsheet_id)
print("--- SQL 스키마 추출 시작 (병렬 처리) ---")
# 내부적으로 _execute_parallel을 호출하여 SCHEMA 시트의 TRUE 항목을 병렬 처리합니다.
schema_sql = await exporter.export_schema()
print(f"{schema_sql}")
with open("output_schema.sql", "w", encoding="utf-8") as f:
f.write(schema_sql)
print("✅ 스키마 추출 완료: output_schema.sql")
print("\n--- 데이터 INSERT 구문 추출 시작 (병렬 처리) ---")
# 내부적으로 _execute_parallel을 호출하여 DATA 시트의 TRUE 항목을 병렬 처리합니다.
data_sql = await exporter.export_data()
print(f"{data_sql}")
with open("output_data.sql", "w", encoding="utf-8") as f:
f.write(data_sql)
print("✅ 데이터 추출 완료: output_data.sql")
# Turn on debugging if you're new to asyncio!
asyncio.run(main(), debug=True)
SpreadsheetExporter.py
import pandas as pd
import numpy as np
import asyncio
from gspread.exceptions import SpreadsheetNotFound, APIError, GSpreadException
class SpreadsheetExporter:
def __init__(self, agc, spreadsheet_id):
"""
:param agc: gspread_asyncio의 AsyncioGspreadClient
:param spreadsheet_id: 메인 SCHEMA 시트 ID
"""
self.agc = agc
self.spreadsheet_id = spreadsheet_id
def _split_column(self, raw_text):
# 1. 원본 문자열 정의
#raw_text = "reactable_type, reactable_id, id"
# 2. 쉼표(,)를 기준으로 분리하여 리스트에 저장 (공백 제거 포함)
column_list = [item.strip() for item in raw_text.split(',')]
# 3. 각 요소를 백틱(`)으로 감싼 후 다시 합치기
result = ", ".join([f"`{item}`" for item in column_list])
# 4. 결과 출력
#print(result)
return result
async def _export_schema_unit(self, row):
"""개별 테이블의 스키마를 생성하는 비동기 단위 작업"""
try:
# 구조 분해 할당 (시트 컬럼명에 맞춰 조정 필요)
table_name = row['TABLE_NAME']
spreadsheet_id = row['SPREADSHEET_ID']
engine = row.get('ENGINE', 'InnoDB')
auto_inc = row.get('AUTO_INCREMENT', '1')
charset = row.get('CHARSET', 'utf8mb4')
collate = row.get('COLLATE', 'utf8mb4_unicode_ci')
ss = await self.agc.open_by_key(spreadsheet_id)
# --- COLUMN, INDEX, FK 시트 데이터를 동시에 가져오기 (병렬) ---
ws_names = ["COLUMN", "INDEX", "FOREIGN_KEY"]
worksheets = await asyncio.gather(*[ss.worksheet(name) for name in ws_names])
data_lists = await asyncio.gather(*[ws.get_all_values() for ws in worksheets])
col_data, idx_data, fk_data = data_lists
all_defs = []
# 1. COLUMN 처리
if len(col_data) > 1:
df = pd.DataFrame(col_data[1:], columns=col_data[0])
#df = df[df.iloc[:, 0] != ""].copy()
"""
# 1. 첫 번째 열이 공백인지 확인 (True/False)
is_empty = (df.iloc[:, 0] == "")
# 2. cummax()를 사용하여 한 번 True가 나오면 그 이후는 계속 True가 되도록 함
# 3. 그 반대(~)인 행들만 남김
df = df[~is_empty.cummax()].copy()
"""
# 2. 첫 번째 열에 글자가 '있는지' 확인 (글자 있으면 True, 비었으면 False)
is_not_empty = (df.iloc[:, 0] != "")
# 3. cumprod()를 사용하여 한 번 False(공백)를 만나는 순간 그 뒤는 무조건 False로 고정
# (중간에 첫 칸이 비면 다음 행에 데이터가 있어도 무시하는 핵심 로직)
df = df[is_not_empty.cumprod().astype(bool)].copy()
#for row in df.itertuples():
#print(f"이름: {row.이름}, 점수: {row.점수}")
type_sql = df.apply(lambda r: f"{r['TYPE']}({r['LENGTH']})" if r.get('LENGTH') else r['TYPE'], axis=1)
null_sql = df['NULL'].apply(lambda x: "NOT NULL" if str(x).upper() in ["NOT", "NO", "FALSE", "0"] else "")
default_sql = df['DEFAULT'].apply(lambda x: f"DEFAULT '{str(x).replace("'", "''")}'" if x else "")
comment_sql = df['COMMENT'].apply(lambda x: f"COMMENT '{str(x)[:1024].replace("'", "''")}'" if x else "")
lines = " `" + df['FIELD'] + "` " + type_sql + " " + null_sql + " " + default_sql + " " + df['EXTRA'] + " " + comment_sql
all_defs.extend(lines.str.strip().tolist())
# 2. INDEX 처리
if len(idx_data) > 1:
df = pd.DataFrame(idx_data[1:], columns=idx_data[0])
#df = df[df.iloc[:, 0] != ""].copy()
# 1. 첫 번째 열이 공백인지 확인 (True/False)
is_empty = (df.iloc[:, 0] == "")
# 2. cummax()를 사용하여 한 번 True가 나오면 그 이후는 계속 True가 되도록 함
# 3. 그 반대(~)인 행들만 남김
df = df[~is_empty.cummax()].copy()
idx_lines = df.apply(lambda r: f"PRIMARY KEY (`{r['COLUMN']}`)" if r['TYPE'].upper() == "PRIMARY"
#else f"{r['TYPE']+' ' if r['TYPE'].upper()!='KEY' else ''}KEY `{r['NAME']}` (`{r['COLUMN']}`)", axis=1)
else f"{r['TYPE']+' ' if r['TYPE'].upper()!='KEY' else ''}KEY `{r['NAME']}` ({self._split_column(r['COLUMN'])})", axis=1)
all_defs.extend(idx_lines.tolist())
# 3. FOREIGN KEY 처리
if len(fk_data) > 1:
df = pd.DataFrame(fk_data[1:], columns=fk_data[0])
#df = df[df.iloc[:, 0] != ""].copy()
# 1. 첫 번째 열이 공백인지 확인 (True/False)
is_empty = (df.iloc[:, 0] == "")
# 2. cummax()를 사용하여 한 번 True가 나오면 그 이후는 계속 True가 되도록 함
# 3. 그 반대(~)인 행들만 남김
df = df[~is_empty.cummax()].copy()
fk_lines = "CONSTRAINT `" + df['NAME'] + "` FOREIGN KEY (`" + df['COLUMN'] + "`) REFERENCES `" + df['REF_TABLE'] + "` (`" + df['REF_COLUMN'] + "`)"
fk_lines += df['ON_UPDATE'].apply(lambda x: f" ON UPDATE {x}" if x else "")
fk_lines += df['ON_DELETE'].apply(lambda x: f" ON DELETE {x}" if x else "")
all_defs.extend(fk_lines.tolist())
return (f"DROP TABLE IF EXISTS `{table_name}`;\n"
f"CREATE TABLE `{table_name}` (\n"
f"{',\n'.join(all_defs)}\n"
f") ENGINE={engine} AUTO_INCREMENT={auto_inc} DEFAULT CHARSET={charset} COLLATE={collate};\n")
except Exception as e:
return f"-- Error in {row.get('TABLE_NAME')}: {e}\n"
async def _export_data_unit(self, row):
"""개별 테이블의 데이터를 INSERT 문으로 생성하는 비동기 단위 작업"""
try:
table_name = row['TABLE_NAME']
ss = await self.agc.open_by_key(row['SPREADSHEET_ID'])
ws = await ss.worksheet("DATA")
data = await ws.get_all_values()
if len(data) <= 1: return ""
df = pd.DataFrame(data[1:], columns=data[0])
#df = df[df.iloc[:, 0] != ""].copy()
# 1. 첫 번째 열이 공백인지 확인 (True/False)
is_empty = (df.iloc[:, 0] == "")
# 2. cummax()를 사용하여 한 번 True가 나오면 그 이후는 계속 True가 되도록 함
# 3. 그 반대(~)인 행들만 남김
df = df[~is_empty.cummax()].copy()
# 데이터 벡터화 처리 (Pandas)
for col in df.columns:
df[col] = df[col].astype(str).str.replace("'", "''")
mask = ~df[col].str.match(r'^-?\d+(\.\d+)?$') # 숫자가 아니면 따옴표
df.loc[mask, col] = "'" + df[col].loc[mask] + "'"
df.loc[df[col] == "''", col] = "NULL"
values_sql = "(" + df.agg(', '.join, axis=1) + ")"
header = ", ".join([f"`{h}`" for h in data[0]])
return f"INSERT INTO `{table_name}` ({header})\nVALUES\n{',\n'.join(values_sql)};\n"
except Exception as e:
return f"-- Data Error in {row.get('TABLE_NAME')}: {e}\n"
async def export_schema(self):
"""메인 실행 함수: 모든 테이블 스키마 병렬 처리"""
return await self._execute_parallel(self._export_schema_unit)
async def export_data(self):
"""메인 실행 함수: 모든 테이블 데이터 병렬 처리"""
return await self._execute_parallel(self._export_data_unit)
async def _execute_parallel(self, unit_func):
"""SCHEMA 시트를 읽고 병렬로 태스크를 실행하는 공통 로직"""
try:
main_ss = await self.agc.open_by_key(self.spreadsheet_id)
schema_ws = await main_ss.worksheet("SCHEMA")
raw_data = await schema_ws.get_all_values()
if len(raw_data) < 2: return "No data found."
df_main = pd.DataFrame(raw_data[1:], columns=raw_data[0])
#df = df[df.iloc[:, 0] != ""].copy()
# 1. 첫 번째 열이 공백인지 확인 (True/False)
is_empty = (df_main.iloc[:, 0] == "")
# 2. cummax()를 사용하여 한 번 True가 나오면 그 이후는 계속 True가 되도록 함
# 3. 그 반대(~)인 행들만 남김
df_main = df_main[~is_empty.cummax()].copy()
# CHK 컬럼이 TRUE인 대상만 선별
active_targets = df_main[df_main.iloc[:, 0].str.upper().str.strip() == "TRUE"]
# --- 핵심: asyncio.gather를 통한 병렬 실행 ---
tasks = [unit_func(row) for _, row in active_targets.iterrows()]
#tasks = [unit_func(row) for row in active_targets.itertuples()]
results = await asyncio.gather(*tasks)
return "\n".join(results)
except Exception as e:
return f"Critical Error: {e}"
"""
⚡ 최적화 포인트 설명
계층적 병렬화 (Double asyncio.gather):
상위 레벨: _execute_parallel에서 여러 테이블(Table A, Table B...)을 동시에 처리합니다.
하위 레벨: _export_schema_unit 내부에서 한 테이블의 COLUMN, INDEX, FK 시트 3개를 동시에 읽어옵니다. (기존 대비 시트 로딩 속도 약 3배 향상)
Pandas 벡터화 연산: for 루프 없이 수천 줄의 INSERT 값이나 DDL 구문을 한 번에 생성합니다.
메모리 효율: df[df.iloc[:, 0] != ""] 코드를 통해 구글 시트 하단의 불필요한 빈 행을 즉시 제거하여 메모리 낭비를 방지합니다.
안정성: asyncio.gather 도중 하나의 시트에서 에러가 발생해도 전체 프로세스가 죽지 않도록 개별 유닛 함수에 try-except를 배치하여 에러 메시지를 SQL 주석(-- Error...) 형태로 반환하게 했습니다.
이제 이 코드를 사용하면 구글 API 할당량(Quota) 내에서 가장 빠른 속도로 SQL 추출이 가능합니다. gspread-asyncio 공식 문서를 참고하여 AsyncioGspreadClient를 넘겨주시면 바로 작동합니다.
"""
"""
"reactable_type, reactable_id, id" 문자열을 , 로 분리후 리스트에 저장한후
"`reactable_type`, `reactable_id`, `id`" 형식으로 화면에 출력해주는 파이썬 코드를 생성해줘
# 1. 원본 문자열 정의
raw_text = "reactable_type, reactable_id, id"
# 2. 쉼표(,)를 기준으로 분리하여 리스트에 저장 (공백 제거 포함)
column_list = [item.strip() for item in raw_text.split(',')]
# 3. 각 요소를 백틱(`)으로 감싼 후 다시 합치기
result = ", ".join([f"`{item}`" for item in column_list])
# 4. 결과 출력
print(result)
"""
app/domain/sheet/__init__.py
# domain/sheet/__init__.py
# 하위 파일들(router, schemas, crud)을 패키지 레벨로 끌어올립니다.
from . import crud, router, schemas
# 외부에서 'from domain.answer import *'를 하거나 가져갈 수 있는 범위를 명시합니다.
__all__ = ["router", "schemas", "crud"]
피드 구독하기:
글 (Atom)
database.py
import os from contextlib import contextmanager # contextmanager 임포트 from typing import Generator from dotenv...
-
bool atob(const char * string) { if (!strcmp(string, "true")) return true; return false; }
-
/// CXXXView.cpp void CXXXView::OnInitialUpdate() { CView::OnInitialUpdate(); // TODO: Add your specialized code here and/or call the ...
-
WxWidgets: http://www.wxwidgets.org/downloads/ MinGW: http://sourceforge.net/projects/mingw/files/ * Microsoft Windows Environment Var...