From 933a7572d526204a7cdfef479b5376062b7d711f Mon Sep 17 00:00:00 2001 From: Kylesoda <249518290+kylesoda@users.noreply.github.com> Date: Tue, 21 Apr 2026 15:00:00 -0500 Subject: [PATCH] feat: integrate azure storage handling in migration pipeline Stream transformer output to an Azure blob container when a job declares ToStorage, with per-job prefix and credential configuration. --- .atl/skill-registry.md | 24 +++ .env.example | 14 +- .gitignore | 3 +- Makefile | 80 ++++++++++ cmd/go_migrate/connect.go | 77 ---------- cmd/go_migrate/log.go | 11 +- cmd/go_migrate/main.go | 28 ++-- cmd/go_migrate/metrics.go | 13 -- cmd/go_migrate/process.go | 44 +++++- config.yaml | 26 +++- go.mod | 8 +- go.sum | 11 +- internal/app/azure/main.go | 87 +++++++++++ internal/app/config/main.go | 43 +++--- internal/app/config/migration.go | 37 +++-- internal/app/etl/extractors/main.go | 101 +++++++++++++ internal/app/etl/extractors/mssql.go | 140 +++++------------- internal/app/etl/extractors/postgres.go | 82 ++++++---- .../app/etl/loaders/{postgres.go => main.go} | 24 +-- internal/app/etl/loaders/types.go | 1 - internal/app/etl/loaders/utils.go | 11 ++ internal/app/etl/table_analyzers/main.go | 15 ++ internal/app/etl/table_analyzers/mssql.go | 6 +- internal/app/etl/table_analyzers/postgres.go | 2 +- internal/app/etl/transformers/mssql.go | 79 +++++++++- internal/app/etl/types.go | 15 +- internal/app/models/main.go | 16 +- 27 files changed, 675 insertions(+), 323 deletions(-) create mode 100644 .atl/skill-registry.md create mode 100644 Makefile delete mode 100644 cmd/go_migrate/connect.go delete mode 100644 cmd/go_migrate/metrics.go create mode 100644 internal/app/azure/main.go create mode 100644 internal/app/etl/extractors/main.go rename internal/app/etl/loaders/{postgres.go => main.go} (81%) delete mode 100644 internal/app/etl/loaders/types.go create mode 100644 internal/app/etl/loaders/utils.go diff --git a/.atl/skill-registry.md b/.atl/skill-registry.md new file mode 100644 index 0000000..2f91391 --- /dev/null +++ b/.atl/skill-registry.md @@ -0,0 +1,24 @@ +# Skill Registry — go-migrate + +Generated: 2026-04-21 + +## Compact Rules + +### Go conventions +- Use existing error wrapping pattern: `fmt.Errorf("context: %w", err)` +- Channel-based pipeline — keep goroutine lifecycle clean (close channels in correct order) +- No comments unless non-obvious WHY; no docstrings +- Prefer named returns only when it aids clarity in short functions +- Use `strings.EqualFold` for case-insensitive column name comparison + +### Project conventions +- Config structs live in `internal/app/config/` +- ETL interfaces live in `internal/app/etl/types.go` +- Transformer implementations in `internal/app/etl/transformers/` +- Azure operations via `internal/app/azure/main.go` +- Per-job transformer creation (not shared) when job has storage config + +## User Skills +| Trigger | Skill | +|---------|-------| +| sdd-* | SDD workflow skills | diff --git a/.env.example b/.env.example index 45ae4a4..803af83 100644 --- a/.env.example +++ b/.env.example @@ -1,2 +1,12 @@ -PG_FROM_DB_URL=postgresql://postgres:password@localhost:5432/db -PG_TO_DB_URL=postgresql://postgres:password@localhost:5432/db \ No newline at end of file +SOURCE_DB_URL=sqlserver://sa:password@localhost:1433?database=master&packet+size=32767&loc=UTC +TARGET_DB_URL=postgresql://postgres:password@localhost:5432/db + +LOG_LEVEL=INFO + +AZ_STORAGE_ENABLED=false +AZ_ACCOUNT_NAME= +AZ_CONTAINER= +AZ_ACCOUNT_KEY= +AZ_USE_HTTPS=true +AZ_SERVICE_URL= +AZ_PREFIX= diff --git a/.gitignore b/.gitignore index 32ba812..9eb588a 100644 --- a/.gitignore +++ b/.gitignore @@ -4,6 +4,7 @@ *.dll *.so *.dylib +bin/ # Test binary, built with `go test -c` *.test @@ -26,5 +27,5 @@ go.work.sum # Editor/IDE # .idea/ -# .vscode/ +.vscode/ .temp diff --git a/Makefile b/Makefile new file mode 100644 index 0000000..394f4e5 --- /dev/null +++ b/Makefile @@ -0,0 +1,80 @@ +.PHONY: build build-linux build-windows build-all clean help + +# Variables +BINARY_NAME=go-migrate +CMD_PATH=./cmd/go_migrate +OUTPUT_DIR=bin +VERSION?=$(shell git describe --tags --always --dirty 2>/dev/null || echo "dev") +BUILD_TIME=$(shell date -u '+%Y-%m-%d_%H:%M:%S') +GIT_COMMIT=$(shell git rev-parse --short HEAD 2>/dev/null || echo "unknown") + +# Flags de compilación +LD_FLAGS=-ldflags="-s -w -X main.Version=$(VERSION) -X main.BuildTime=$(BUILD_TIME) -X main.GitCommit=$(GIT_COMMIT)" + +# Default: compilar para el SO actual +build: build-$(OS) + +ifeq ($(OS),Windows_NT) +build-native: build-windows +else +build-native: build-linux +endif + +# Compilar para Linux (sin CGO para máxima compatibilidad) +build-linux: + @echo "Compilando para Linux..." + @mkdir -p $(OUTPUT_DIR) + CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build \ + $(LD_FLAGS) \ + -o $(OUTPUT_DIR)/$(BINARY_NAME)-linux-amd64 \ + $(CMD_PATH) + @echo "Binario creado: $(OUTPUT_DIR)/$(BINARY_NAME)-linux-amd64" + +# Compilar para Windows +build-windows: + @echo "Compilando para Windows..." + @mkdir -p $(OUTPUT_DIR) + CGO_ENABLED=0 GOOS=windows GOARCH=amd64 go build \ + $(LD_FLAGS) \ + -o $(OUTPUT_DIR)/$(BINARY_NAME)-windows-amd64.exe \ + $(CMD_PATH) + @echo "Binario creado: $(OUTPUT_DIR)/$(BINARY_NAME)-windows-amd64.exe" + +# Compilar para ambas plataformas +build-all: build-linux build-windows + @echo "" + @echo "Binarios compilados:" + @ls -lh $(OUTPUT_DIR)/$(BINARY_NAME)* + +# Compilar para Linux arm64 (opcional, para Raspberry Pi, etc.) +build-linux-arm64: + @echo "Compilando para Linux ARM64..." + @mkdir -p $(OUTPUT_DIR) + CGO_ENABLED=0 GOOS=linux GOARCH=arm64 go build \ + $(LD_FLAGS) \ + -o $(OUTPUT_DIR)/$(BINARY_NAME)-linux-arm64 \ + $(CMD_PATH) + @echo "Binario creado: $(OUTPUT_DIR)/$(BINARY_NAME)-linux-arm64" + +# Limpiar binarios +clean: + @echo "Limpiando binarios..." + @rm -rf $(OUTPUT_DIR) + @echo "Limpieza completada" + +# Ayuda +help: + @echo "Comandos disponibles:" + @echo "" + @echo " make build - Compilar para el SO actual (Linux/Windows)" + @echo " make build-linux - Compilar para Linux x86_64" + @echo " make build-windows - Compilar para Windows x86_64" + @echo " make build-linux-arm64 - Compilar para Linux ARM64 (opcional)" + @echo " make build-all - Compilar para Linux y Windows" + @echo " make clean - Eliminar binarios compilados" + @echo " make help - Mostrar esta ayuda" + @echo "" + @echo "Ejemplos de uso:" + @echo " make build-all # Crear binarios para ambas plataformas" + @echo " make build-linux OS= # Crear solo para Linux" + @echo "" diff --git a/cmd/go_migrate/connect.go b/cmd/go_migrate/connect.go deleted file mode 100644 index 40e417f..0000000 --- a/cmd/go_migrate/connect.go +++ /dev/null @@ -1,77 +0,0 @@ -package main - -import ( - "context" - "database/sql" - "errors" - "fmt" - "sync" - "time" - - "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/config" - "github.com/jackc/pgx/v5/pgxpool" - _ "github.com/microsoft/go-mssqldb" - log "github.com/sirupsen/logrus" -) - -func connectToSqlServer() (*sql.DB, error) { - db, err := sql.Open("sqlserver", config.App.SourceDbUrl) - if err != nil { - return nil, fmt.Errorf("Unable to connect to sqlserver: %w", err) - } - - ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second) - defer cancel() - - if err := db.PingContext(ctx); err != nil { - return nil, fmt.Errorf("Unable to ping sqlserver: %w", err) - } - - return db, nil -} - -func connectToPostgres() (*pgxpool.Pool, error) { - ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second) - defer cancel() - - pool, err := pgxpool.New(ctx, config.App.TargetDbUrl) - if err != nil { - return nil, fmt.Errorf("Unable to connect to postgres: %w", err) - } - - if err := pool.Ping(ctx); err != nil { - pool.Close() - return nil, fmt.Errorf("Unable to ping postgres: %w", err) - } - - return pool, nil -} - -func connectToDatabases() (*sql.DB, *pgxpool.Pool, error) { - var sourceDbErr, targetDbErr error - var sourceDb *sql.DB - var targetDb *pgxpool.Pool - var wg sync.WaitGroup - - wg.Go(func() { - sourceDb, sourceDbErr = connectToSqlServer() - if sourceDbErr != nil { - log.Error("Unable to connect to source db: ", sourceDbErr) - } - }) - - wg.Go(func() { - targetDb, targetDbErr = connectToPostgres() - if targetDbErr != nil { - log.Error("Unable to connect to target db: ", targetDbErr) - } - }) - - wg.Wait() - - if sourceDbErr != nil || targetDbErr != nil { - return nil, nil, errors.New("Unable to connect to databases") - } - - return sourceDb, targetDb, nil -} diff --git a/cmd/go_migrate/log.go b/cmd/go_migrate/log.go index a7afee3..29c1934 100644 --- a/cmd/go_migrate/log.go +++ b/cmd/go_migrate/log.go @@ -3,6 +3,7 @@ package main import ( "time" + "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/config" log "github.com/sirupsen/logrus" ) @@ -13,5 +14,13 @@ func configureLog() { DisableSorting: false, PadLevelText: true, }) - log.SetLevel(log.DebugLevel) + + logLevelEnv := config.App.LogLevel + logLevel, err := log.ParseLevel(logLevelEnv) + if err != nil { + log.Warnf("Nivel de log inválido '%s', usando INFO por defecto", logLevelEnv) + logLevel = log.InfoLevel + } + + log.SetLevel(logLevel) } diff --git a/cmd/go_migrate/main.go b/cmd/go_migrate/main.go index e57fe3f..9ed7816 100644 --- a/cmd/go_migrate/main.go +++ b/cmd/go_migrate/main.go @@ -5,12 +5,13 @@ import ( "sync" "time" + "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/azure" "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/config" - "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/db-wrapper" + dbwrapper "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/db-wrapper" "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/etl/extractors" "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/etl/loaders" "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/etl/table_analyzers" - "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/etl/transformers" + "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/models" log "github.com/sirupsen/logrus" "golang.org/x/sync/errgroup" ) @@ -95,10 +96,10 @@ func processMigrationJobs( targetDb dbwrapper.DbWrapper, jobs []config.Job, maxParallelWorkers int, -) []JobResult { +) []models.JobResult { if len(jobs) == 0 { log.Info("No migration jobs configured") - return []JobResult{} + return []models.JobResult{} } if maxParallelWorkers <= 0 { @@ -111,15 +112,23 @@ func processMigrationJobs( log.Infof("Starting migration with %d parallel worker(s)", maxParallelWorkers) - chJobResults := make(chan JobResult, len(jobs)) + chJobResults := make(chan models.JobResult, len(jobs)) chJobs := make(chan config.Job, len(jobs)) var wgJobs sync.WaitGroup sourceTableAnalyzer := table_analyzers.NewMssqlTableAnalyzer(sourceDb) targetTableAnalyzer := table_analyzers.NewPostgresTableAnalyzer(targetDb) extractor := extractors.NewMssqlExtractor(sourceDb) - transformer := transformers.NewMssqlTransformer() - loader := loaders.NewPostgresLoader(targetDb) + loader := loaders.NewGenericLoader(targetDb) + + var azureClient *azure.Client + if config.App.AzureStorage.Enabled { + var err error + azureClient, err = azure.NewClient(config.App.AzureStorage) + if err != nil { + log.Fatalf("Failed to create Azure storage client: %v", err) + } + } for i := range maxParallelWorkers { wgJobs.Go(func() { @@ -131,9 +140,10 @@ func processMigrationJobs( sourceTableAnalyzer, targetTableAnalyzer, extractor, - transformer, + azureClient, loader, job, + targetDb.GetDialect(), ) chJobResults <- res @@ -151,7 +161,7 @@ func processMigrationJobs( close(chJobResults) }() - var finalResults []JobResult + var finalResults []models.JobResult for res := range chJobResults { finalResults = append(finalResults, res) } diff --git a/cmd/go_migrate/metrics.go b/cmd/go_migrate/metrics.go deleted file mode 100644 index b540c8c..0000000 --- a/cmd/go_migrate/metrics.go +++ /dev/null @@ -1,13 +0,0 @@ -package main - -import "time" - -type JobResult struct { - JobName string - StartTime time.Time - Duration time.Duration - RowsRead int64 - RowsLoaded int64 - RowsFailed int64 - Error error -} diff --git a/cmd/go_migrate/process.go b/cmd/go_migrate/process.go index cfc678b..b87a1f7 100644 --- a/cmd/go_migrate/process.go +++ b/cmd/go_migrate/process.go @@ -2,34 +2,54 @@ package main import ( "context" + "fmt" "sync" "sync/atomic" "time" + "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/azure" "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/config" "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/custom_errors" dbwrapper "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/db-wrapper" "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/etl" + "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/etl/extractors" "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/etl/table_analyzers" + "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/etl/transformers" "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/models" log "github.com/sirupsen/logrus" "golang.org/x/sync/errgroup" ) +func buildTruncateQuery(targetDbType, schema, table, truncateMethod string) string { + if truncateMethod == "DELETE" { + if targetDbType == "postgres" { + return fmt.Sprintf(`DELETE FROM "%s"."%s"`, schema, table) + } + return fmt.Sprintf(`DELETE FROM [%s].[%s]`, schema, table) + } + + if targetDbType == "postgres" { + return fmt.Sprintf(`TRUNCATE TABLE "%s"."%s"`, schema, table) + } + return fmt.Sprintf(`TRUNCATE TABLE [%s].[%s]`, schema, table) +} + func processMigrationJob( ctx context.Context, targetDbWrapper dbwrapper.DbWrapper, sourceTableAnalyzer etl.TableAnalyzer, targetTableAnalyzer etl.TableAnalyzer, extractor etl.Extractor, - transformer etl.Transformer, + azureClient *azure.Client, loader etl.Loader, job config.Job, -) JobResult { + targetDbType string, +) models.JobResult { + transformer := transformers.NewMssqlTransformer(job.ToStorage, job.SourceTable, azureClient) localCtx, cancel := context.WithCancel(ctx) defer cancel() - result := JobResult{ + result := models.JobResult{ JobName: job.Name, StartTime: time.Now(), } @@ -65,7 +85,13 @@ func processMigrationJob( return result } - for _, query := range job.PreSQL { + preSqlQueries := job.TargetTable.PreSQL + if job.TruncateTarget { + truncateQuery := buildTruncateQuery(targetDbType, job.TargetTable.Schema, job.TargetTable.Table, job.TruncateMethod) + preSqlQueries = append([]string{truncateQuery}, job.TargetTable.PreSQL...) + } + + for _, query := range preSqlQueries { if _, err := targetDbWrapper.Exec(localCtx, query); err != nil { result.Error = err return result @@ -78,6 +104,7 @@ func processMigrationJob( job.SourceTable.TableInfo, job.SourceTable.PrimaryKey, job.RowsPerPartition, + job.Range, ) if err != nil { log.Error("Unexpected error calculating batch ranges: ", err) @@ -128,8 +155,9 @@ func processMigrationJob( for range maxExtractors { wgExtractors.Go(func() { - extractor.Exec( + extractors.Consume( localCtx, + extractor, job.SourceTable, sourceColTypes, job.BatchSize, @@ -213,7 +241,7 @@ func processMigrationJob( cancel() }() - for _, query := range job.PostSQL { + for _, query := range job.TargetTable.PostSQL { if _, err := targetDbWrapper.Exec(localCtx, query); err != nil { result.Error = err return result @@ -233,5 +261,9 @@ func processMigrationJob( result.RowsLoaded = atomic.LoadInt64(&rowsLoaded) result.RowsFailed = atomic.LoadInt64(&rowsFailed) + if result.RowsRead != result.RowsLoaded { + result.Error = fmt.Errorf("Row count mismatch: extracted %d rows but loaded %d rows (failed: %d)", result.RowsRead, result.RowsLoaded, result.RowsFailed) + } + return result } diff --git a/config.yaml b/config.yaml index af68f11..93548ef 100644 --- a/config.yaml +++ b/config.yaml @@ -28,8 +28,8 @@ jobs: target: schema: demo table: users - pre_sql: - - 'SELECT 1' + pre_sql: + - 'SELECT 1' range: min: 1000000 max: 2000000 @@ -45,7 +45,21 @@ jobs: target: schema: analytics table: events - pre_sql: - - 'SELECT 1' - post_sql: - - "SELECT 1" + pre_sql: + - 'SELECT 1' + post_sql: + - "SELECT 1" + + - name: storage_attachments + source: + schema: analytics + table: attachments + primary_key: id + target: + schema: analytics + table: attachments + to_storage: + columns: + - source: DATA + target: FILE_URL + mode: REFERENCE_ONLY # REFERENCE_ONLY | DUPLICATE_WITH_REF diff --git a/go.mod b/go.mod index 96170be..f666b7f 100644 --- a/go.mod +++ b/go.mod @@ -1,8 +1,10 @@ module git.ksdemosapps.com/kylesoda/go-migrate -go 1.25.7 +go 1.26 require ( + github.com/Azure/azure-sdk-for-go/sdk/storage/azblob v1.6.4 + github.com/caarlos0/env/v11 v11.4.0 github.com/gaspardle/go-mssqlclrgeo v0.0.0-20160129143314-97ceabf987a4 github.com/google/uuid v1.6.0 github.com/jackc/pgx/v5 v5.9.1 @@ -15,15 +17,17 @@ require ( ) require ( + github.com/Azure/azure-sdk-for-go/sdk/azcore v1.21.0 // indirect + github.com/Azure/azure-sdk-for-go/sdk/internal v1.11.2 // indirect github.com/golang-sql/civil v0.0.0-20220223132316-b832511892a9 // indirect github.com/golang-sql/sqlexp v0.1.0 // indirect github.com/jackc/pgpassfile v1.0.0 // indirect github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect github.com/jackc/puddle/v2 v2.2.2 // indirect - github.com/kr/text v0.2.0 // indirect github.com/rogpeppe/go-internal v1.14.1 // indirect github.com/shopspring/decimal v1.4.0 // indirect golang.org/x/crypto v0.48.0 // indirect + golang.org/x/net v0.51.0 // indirect golang.org/x/sys v0.41.0 // indirect golang.org/x/text v0.34.0 // indirect ) diff --git a/go.sum b/go.sum index 7fecac7..16550d7 100644 --- a/go.sum +++ b/go.sum @@ -4,10 +4,14 @@ github.com/Azure/azure-sdk-for-go/sdk/azidentity v1.13.1 h1:Hk5QBxZQC1jb2Fwj6mpz github.com/Azure/azure-sdk-for-go/sdk/azidentity v1.13.1/go.mod h1:IYus9qsFobWIc2YVwe/WPjcnyCkPKtnHAqUYeebc8z0= github.com/Azure/azure-sdk-for-go/sdk/internal v1.11.2 h1:9iefClla7iYpfYWdzPCRDozdmndjTm8DXdpCzPajMgA= github.com/Azure/azure-sdk-for-go/sdk/internal v1.11.2/go.mod h1:XtLgD3ZD34DAaVIIAyG3objl5DynM3CQ/vMcbBNJZGI= +github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/storage/armstorage v1.8.1 h1:/Zt+cDPnpC3OVDm/JKLOs7M2DKmLRIIp3XIx9pHHiig= +github.com/Azure/azure-sdk-for-go/sdk/resourcemanager/storage/armstorage v1.8.1/go.mod h1:Ng3urmn6dYe8gnbCMoHHVl5APYz2txho3koEkV2o2HA= github.com/Azure/azure-sdk-for-go/sdk/security/keyvault/azkeys v1.4.0 h1:E4MgwLBGeVB5f2MdcIVD3ELVAWpr+WD6MUe1i+tM/PA= github.com/Azure/azure-sdk-for-go/sdk/security/keyvault/azkeys v1.4.0/go.mod h1:Y2b/1clN4zsAoUd/pgNAQHjLDnTis/6ROkUfyob6psM= github.com/Azure/azure-sdk-for-go/sdk/security/keyvault/internal v1.2.0 h1:nCYfgcSyHZXJI8J0IWE5MsCGlb2xp9fJiXyxWgmOFg4= github.com/Azure/azure-sdk-for-go/sdk/security/keyvault/internal v1.2.0/go.mod h1:ucUjca2JtSZboY8IoUqyQyuuXvwbMBVwFOm0vdQPNhA= +github.com/Azure/azure-sdk-for-go/sdk/storage/azblob v1.6.4 h1:jWQK1GI+LeGGUKBADtcH2rRqPxYB1Ljwms5gFA2LqrM= +github.com/Azure/azure-sdk-for-go/sdk/storage/azblob v1.6.4/go.mod h1:8mwH4klAm9DUgR2EEHyEEAQlRDvLPyg5fQry3y+cDew= github.com/AzureAD/microsoft-authentication-library-for-go v1.6.0 h1:XRzhVemXdgvJqCH0sFfrBUTnUJSBrBf7++ypk+twtRs= github.com/AzureAD/microsoft-authentication-library-for-go v1.6.0/go.mod h1:HKpQxkWaGLJ+D/5H8QRpyQXA1eKjxkFlOMwck5+33Jk= github.com/DATA-DOG/go-sqlmock v1.5.2 h1:OcvFkGmslmlZibjAjaHm3L//6LiuBgolP7OputlJIzU= @@ -16,7 +20,8 @@ github.com/alecthomas/assert/v2 v2.10.0 h1:jjRCHsj6hBJhkmhznrCzoNpbA3zqy0fYiUcYZ github.com/alecthomas/assert/v2 v2.10.0/go.mod h1:Bze95FyfUr7x34QZrjL+XP+0qgp/zg8yS+TtBj1WA3k= github.com/alecthomas/repr v0.4.0 h1:GhI2A8MACjfegCPVq9f1FLvIBS+DrQ2KQBFZP1iFzXc= github.com/alecthomas/repr v0.4.0/go.mod h1:Fr0507jx4eOXV7AlPV6AVZLYrLIuIeSOWtW57eE/O/4= -github.com/creack/pty v1.1.9/go.mod h1:oKZEueFk5CKHvIhNR5MUki03XCEU+Q6VDXinZuGJ33E= +github.com/caarlos0/env/v11 v11.4.0 h1:Kcb6t5kIIr4XkoQC9AF2j+8E1Jsrl3Wz/hhm1LtoGAc= +github.com/caarlos0/env/v11 v11.4.0/go.mod h1:qupehSf/Y0TUTsxKywqRt/vJjN5nz6vauiYEUUr8P4U= github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= @@ -42,8 +47,8 @@ github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= github.com/joho/godotenv v1.5.1 h1:7eLL/+HRGLY0ldzfGMeQkb7vMd0as4CfYvUVzLqw0N0= github.com/joho/godotenv v1.5.1/go.mod h1:f4LDr5Voq0i2e/R5DDNOoa2zzDfwtkZa6DnEwAbqwq4= -github.com/kr/pretty v0.3.0 h1:WgNl7dwNpEZ6jJ9k1snq4pZsg7DOEN8hP9Xw0Tsjwk0= -github.com/kr/pretty v0.3.0/go.mod h1:640gp4NfQd8pI5XOwp5fnNeVWj67G7CFk/SaSQn7NBk= +github.com/kr/pretty v0.3.1 h1:flRD4NNwYAUpkphVc1HcthR4KEIFJ65n8Mw5qdRn3LE= +github.com/kr/pretty v0.3.1/go.mod h1:hoEshYVHaxMs3cyo3Yncou5ZscifuDolrwPKZanG3xk= github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY= github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE= github.com/kylelemons/godebug v1.1.0 h1:RPNrshWIDI6G2gRW9EHilWtl7Z6Sb1BR0xunSBf0SNc= diff --git a/internal/app/azure/main.go b/internal/app/azure/main.go new file mode 100644 index 0000000..7c08bef --- /dev/null +++ b/internal/app/azure/main.go @@ -0,0 +1,87 @@ +package azure + +import ( + "context" + "errors" + "fmt" + "net/url" + + "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/config" + "github.com/Azure/azure-sdk-for-go/sdk/storage/azblob" +) + +var ( + ErrInvalidConnectionString = errors.New("invalid connection string") + ErrContainerNotFound = errors.New("container not found") + ErrBlobNotFound = errors.New("blob not found") + ErrInvalidInput = errors.New("invalid input parameters") +) + +type Client struct { + client *azblob.Client + azureStorageConfig config.AzureStorageConfig +} + +func NewClient(azureStorageConfig config.AzureStorageConfig) (*Client, error) { + protocol := "https" + if !azureStorageConfig.UseHTTPS { + protocol = "http" + } + + blobEndpoint, _ := url.JoinPath(azureStorageConfig.ServiceURL, azureStorageConfig.AccountName) + connStr := fmt.Sprintf("DefaultEndpointsProtocol=%s;AccountName=%s;AccountKey=%s;BlobEndpoint=%s;", + protocol, azureStorageConfig.AccountName, azureStorageConfig.AccountKey, blobEndpoint) + + client, err := azblob.NewClientFromConnectionString(connStr, nil) + if err != nil { + return nil, fmt.Errorf("creating azure storage client: %w", err) + } + + return &Client{ + client: client, + azureStorageConfig: azureStorageConfig, + }, nil +} + +func (c *Client) CreateContainer(ctx context.Context, containerName string) error { + if containerName == "" { + return ErrInvalidInput + } + + _, err := c.client.CreateContainer(ctx, containerName, nil) + if err != nil { + return fmt.Errorf("creating container %s: %w", containerName, err) + } + return nil +} + +func (c *Client) UploadBuffer(ctx context.Context, containerName, blobPath string, buffer []byte) error { + if containerName == "" || blobPath == "" || buffer == nil { + return ErrInvalidInput + } + + _, err := c.client.UploadBuffer(ctx, containerName, blobPath, buffer, nil) + if err != nil { + return fmt.Errorf("uploading blob %s: %w", blobPath, err) + } + return nil +} + +func (c *Client) UploadAndGetURL(ctx context.Context, blobPath string, buffer []byte) (string, error) { + if blobPath == "" || buffer == nil { + return "", ErrInvalidInput + } + + fullPath := blobPath + if c.azureStorageConfig.Prefix != "" { + fullPath, _ = url.JoinPath(c.azureStorageConfig.Prefix, blobPath) + } + + if err := c.UploadBuffer(ctx, c.azureStorageConfig.Container, fullPath, buffer); err != nil { + return "", err + } + + blobEndpoint, _ := url.JoinPath(c.azureStorageConfig.ServiceURL, c.azureStorageConfig.AccountName) + blobURL, _ := url.JoinPath(blobEndpoint, c.azureStorageConfig.Container, fullPath) + return blobURL, nil +} diff --git a/internal/app/config/main.go b/internal/app/config/main.go index 94e751d..1ff1473 100644 --- a/internal/app/config/main.go +++ b/internal/app/config/main.go @@ -1,41 +1,40 @@ package config import ( - "os" - + "github.com/caarlos0/env/v11" "github.com/joho/godotenv" log "github.com/sirupsen/logrus" ) -type appConfig struct { - SourceDbUrl string - TargetDbUrl string +type AzureStorageConfig struct { + AccountName string `env:"AZ_ACCOUNT_NAME"` + Container string `env:"AZ_CONTAINER"` + AccountKey string `env:"AZ_ACCOUNT_KEY"` + UseHTTPS bool `env:"AZ_USE_HTTPS" envDefault:"true"` + ServiceURL string `env:"AZ_SERVICE_URL"` + Prefix string `env:"AZ_PREFIX"` + Enabled bool `env:"AZ_STORAGE_ENABLED"` } -func loadEnv() { - err := godotenv.Load() - if err != nil { - log.Warn("Warning: could not load .env file") - } +type appConfig struct { + SourceDbUrl string `env:"SOURCE_DB_URL,required"` + TargetDbUrl string `env:"TARGET_DB_URL,required"` + LogLevel string `env:"LOG_LEVEL" envDefault:"INFO"` + AzureStorage AzureStorageConfig } func getAppConfig() appConfig { - loadEnv() - - sourceDbUrl := os.Getenv("SOURCE_DB_URL") - if sourceDbUrl == "" { - log.Fatal("SOURCE_DB_URL environment variable not set") + if err := godotenv.Load(); err != nil { + log.Warn("Could not load .env file") } - targetDbUrl := os.Getenv("TARGET_DB_URL") - if targetDbUrl == "" { - log.Fatal("TARGET_DB_URL environment variable not set") + cfg := appConfig{} + + if err := env.Parse(&cfg); err != nil { + log.Fatalf("Error al cargar variables de entorno: %v", err) } - return appConfig{ - SourceDbUrl: sourceDbUrl, - TargetDbUrl: targetDbUrl, - } + return cfg } var App appConfig = getAppConfig() diff --git a/internal/app/config/migration.go b/internal/app/config/migration.go index cfd4d06..bfbae35 100644 --- a/internal/app/config/migration.go +++ b/internal/app/config/migration.go @@ -14,6 +14,16 @@ type RetryConfig struct { MaxJitterMs int `yaml:"max_jitter_ms"` } +type ToStorageColumnConfig struct { + Source string `yaml:"source"` + Target string `yaml:"target"` + Mode string `yaml:"mode"` +} + +type ToStorageConfig struct { + Columns []ToStorageColumnConfig `yaml:"columns"` +} + type JobConfig struct { MaxExtractors int `yaml:"max_extractors"` MaxLoaders int `yaml:"max_loaders"` @@ -26,6 +36,7 @@ type JobConfig struct { MaxChunkErrors int `yaml:"max_chunk_errors"` Retry RetryConfig `yaml:"retry"` RowsPerPartition int64 + ToStorage ToStorageConfig `yaml:"to_storage"` } type TableInfo struct { @@ -33,29 +44,31 @@ type TableInfo struct { Table string `yaml:"table"` } -type TargetTableInfo struct { - TableInfo `yaml:",inline"` -} - type SourceTableInfo struct { TableInfo `yaml:",inline"` PrimaryKey string `yaml:"primary_key"` } +type TargetTableInfo struct { + TableInfo `yaml:",inline"` + PreSQL []string `yaml:"pre_sql"` + PostSQL []string `yaml:"post_sql"` +} + +type RangeConfig struct { + Min int64 `yaml:"min"` + Max int64 `yaml:"max"` + IsMinInclusive bool `yaml:"is_min_inclusive"` + IsMaxInclusive bool `yaml:"is_max_inclusive"` +} + type Job struct { Name string `yaml:"name"` Enabled bool `yaml:"enabled"` SourceTable SourceTableInfo `yaml:"source"` TargetTable TargetTableInfo `yaml:"target"` - PreSQL []string `yaml:"pre_sql"` - PostSQL []string `yaml:"post_sql"` JobConfig `yaml:",inline"` - Range struct { - Min int64 `yaml:"min"` - Max int64 `yaml:"max"` - IsMinInclusive bool `yaml:"is_min_inclusive"` - IsMaxInclusive bool `yaml:"is_max_inclusive"` - } + Range RangeConfig `yaml:"range"` } type MigrationConfig struct { diff --git a/internal/app/etl/extractors/main.go b/internal/app/etl/extractors/main.go new file mode 100644 index 0000000..d3137ac --- /dev/null +++ b/internal/app/etl/extractors/main.go @@ -0,0 +1,101 @@ +package extractors + +import ( + "context" + "errors" + "slices" + "strings" + "sync" + "sync/atomic" + + "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/config" + "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/custom_errors" + "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/etl" + "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/models" +) + +func Consume( + ctx context.Context, + extractor etl.Extractor, + tableInfo config.SourceTableInfo, + columns []models.ColumnType, + batchSize int, + chPartitionsIn <-chan models.Partition, + chBatchesOut chan<- models.Batch, + chErrorsOut chan<- custom_errors.ExtractorError, + chJobErrorsOut chan<- custom_errors.JobError, + wgActivePartitions *sync.WaitGroup, + rowsRead *int64, +) { + indexPrimaryKey := slices.IndexFunc(columns, func(col models.ColumnType) bool { + return strings.EqualFold(col.Name(), tableInfo.PrimaryKey) + }) + + if indexPrimaryKey == -1 { + select { + case <-ctx.Done(): + return + case chJobErrorsOut <- custom_errors.JobError{ + ShouldCancelJob: true, + Msg: "Primary key not found in provided columns", + }: + } + + return + } + + for { + if ctx.Err() != nil { + return + } + + select { + case <-ctx.Done(): + return + case partition, ok := <-chPartitionsIn: + if !ok { + return + } + + rowsReadResult, err := extractor.Exec( + ctx, + tableInfo, + columns, + batchSize, + partition, + indexPrimaryKey, + chBatchesOut, + ) + + if rowsReadResult > 0 { + atomic.AddInt64(rowsRead, int64(rowsReadResult)) + } + + if err != nil { + if exError, ok := errors.AsType[*custom_errors.ExtractorError](err); ok { + select { + case <-ctx.Done(): + return + case chErrorsOut <- *exError: + } + } else if jobError, ok := errors.AsType[*custom_errors.JobError](err); ok { + select { + case <-ctx.Done(): + return + case chJobErrorsOut <- *jobError: + } + } else { + select { + case <-ctx.Done(): + return + case chErrorsOut <- custom_errors.ExtractorError{Partition: partition, Msg: err.Error()}: + } + } + + continue + } + + wgActivePartitions.Done() + } + } +} diff --git a/internal/app/etl/extractors/mssql.go b/internal/app/etl/extractors/mssql.go index d34447a..e2b1300 100644 --- a/internal/app/etl/extractors/mssql.go +++ b/internal/app/etl/extractors/mssql.go @@ -5,10 +5,7 @@ import ( "database/sql" "errors" "fmt" - "slices" "strings" - "sync" - "sync/atomic" "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/config" "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/convert" @@ -32,6 +29,9 @@ func buildExtractQueryMssql( columns []models.ColumnType, includeRange bool, isMinInclusive bool, + isMaxInclusive bool, + hasMin bool, + hasMax bool, ) string { var sbQuery strings.Builder @@ -55,15 +55,32 @@ func buildExtractQueryMssql( fmt.Fprintf(&sbQuery, " FROM [%s].[%s]", tableInfo.Schema, tableInfo.Table) - if includeRange { - fmt.Fprintf(&sbQuery, " WHERE [%s]", tableInfo.PrimaryKey) - if isMinInclusive { - sbQuery.WriteString(" >=") - } else { - sbQuery.WriteString(" >") + if includeRange && (hasMin || hasMax) { + sbQuery.WriteString(" WHERE ") + + if hasMin { + fmt.Fprintf(&sbQuery, "[%s]", tableInfo.PrimaryKey) + if isMinInclusive { + sbQuery.WriteString(" >=") + } else { + sbQuery.WriteString(" >") + } + sbQuery.WriteString(" @min") } - fmt.Fprintf(&sbQuery, " @min AND [%s] <= @max", tableInfo.PrimaryKey) + if hasMin && hasMax { + sbQuery.WriteString(" AND ") + } + + if hasMax { + fmt.Fprintf(&sbQuery, "[%s]", tableInfo.PrimaryKey) + if isMaxInclusive { + sbQuery.WriteString(" <=") + } else { + sbQuery.WriteString(" <") + } + sbQuery.WriteString(" @max") + } } fmt.Fprintf(&sbQuery, " ORDER BY [%s] ASC", tableInfo.PrimaryKey) @@ -99,7 +116,7 @@ func errorFromLastRow( } } -func (mssqlEx *MssqlExtractor) ProcessPartition( +func (mssqlEx *MssqlExtractor) Exec( ctx context.Context, tableInfo config.SourceTableInfo, columns []models.ColumnType, @@ -108,14 +125,16 @@ func (mssqlEx *MssqlExtractor) ProcessPartition( indexPrimaryKey int, chBatchesOut chan<- models.Batch, ) (int, error) { - query := buildExtractQueryMssql(tableInfo, columns, partition.HasRange, partition.Range.IsMinInclusive) + hasMin := partition.HasRange && partition.Range.Min > 0 + hasMax := partition.HasRange && partition.Range.Max > 0 + query := buildExtractQueryMssql(tableInfo, columns, partition.HasRange, partition.Range.IsMinInclusive, partition.Range.IsMaxInclusive, hasMin, hasMax) var queryArgs []any - if partition.HasRange { - queryArgs = append(queryArgs, - sql.Named("min", partition.Range.Min), - sql.Named("max", partition.Range.Max), - ) + if hasMin { + queryArgs = append(queryArgs, sql.Named("min", partition.Range.Min)) + } + if hasMax { + queryArgs = append(queryArgs, sql.Named("max", partition.Range.Max)) } rowsRead := 0 @@ -188,90 +207,3 @@ func (mssqlEx *MssqlExtractor) ProcessPartition( return rowsRead, nil } - -func (mssqlEx *MssqlExtractor) Exec( - ctx context.Context, - tableInfo config.SourceTableInfo, - columns []models.ColumnType, - batchSize int, - chPartitionsIn <-chan models.Partition, - chBatchesOut chan<- models.Batch, - chErrorsOut chan<- custom_errors.ExtractorError, - chJobErrorsOut chan<- custom_errors.JobError, - wgActivePartitions *sync.WaitGroup, - rowsRead *int64, -) { - indexPrimaryKey := slices.IndexFunc(columns, func(col models.ColumnType) bool { - return strings.EqualFold(col.Name(), tableInfo.PrimaryKey) - }) - - if indexPrimaryKey == -1 { - select { - case <-ctx.Done(): - return - case chJobErrorsOut <- custom_errors.JobError{ - ShouldCancelJob: true, - Msg: "Primary key not found in provided columns", - }: - } - - return - } - - for { - if ctx.Err() != nil { - return - } - - select { - case <-ctx.Done(): - return - case partition, ok := <-chPartitionsIn: - if !ok { - return - } - - rowsReadResult, err := mssqlEx.ProcessPartition( - ctx, - tableInfo, - columns, - batchSize, - partition, - indexPrimaryKey, - chBatchesOut, - ) - - if rowsReadResult > 0 { - atomic.AddInt64(rowsRead, int64(rowsReadResult)) - } - - if err != nil { - var exError *custom_errors.ExtractorError - var jobError *custom_errors.JobError - if errors.As(err, &exError) { - select { - case <-ctx.Done(): - return - case chErrorsOut <- *exError: - } - } else if errors.As(err, &jobError) { - select { - case <-ctx.Done(): - return - case chJobErrorsOut <- *jobError: - } - } else { - select { - case <-ctx.Done(): - return - case chErrorsOut <- custom_errors.ExtractorError{Partition: partition, Msg: err.Error()}: - } - } - - continue - } - - wgActivePartitions.Done() - } - } -} diff --git a/internal/app/etl/extractors/postgres.go b/internal/app/etl/extractors/postgres.go index 8d494d3..964f940 100644 --- a/internal/app/etl/extractors/postgres.go +++ b/internal/app/etl/extractors/postgres.go @@ -2,10 +2,8 @@ package extractors import ( "context" - "errors" "fmt" "strings" - "sync" "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/config" "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/custom_errors" @@ -23,7 +21,15 @@ func NewPostgresExtractor(db dbwrapper.DbWrapper) etl.Extractor { return &PostgresExtractor{db: db} } -func buildExtractQueryPostgres(sourceDbInfo config.SourceTableInfo, columns []models.ColumnType) string { +func buildExtractQueryPostgres( + sourceDbInfo config.SourceTableInfo, + columns []models.ColumnType, + includeRange bool, + isMinInclusive bool, + isMaxInclusive bool, + hasMin bool, + hasMax bool, +) string { var sbColumns strings.Builder if len(columns) == 0 { @@ -48,10 +54,44 @@ func buildExtractQueryPostgres(sourceDbInfo config.SourceTableInfo, columns []mo } } - return fmt.Sprintf(`SELECT %s FROM "%s"."%s" ORDER BY "%s" ASC`, sbColumns.String(), sourceDbInfo.Schema, sourceDbInfo.Table, sourceDbInfo.PrimaryKey) + query := fmt.Sprintf(`SELECT %s FROM "%s"."%s"`, sbColumns.String(), sourceDbInfo.Schema, sourceDbInfo.Table) + + if includeRange && (hasMin || hasMax) { + query += " WHERE " + paramIdx := 1 + + if hasMin { + query += fmt.Sprintf(`"%s"`, sourceDbInfo.PrimaryKey) + if isMinInclusive { + query += " >=" + } else { + query += " >" + } + query += fmt.Sprintf(" $%d", paramIdx) + paramIdx++ + } + + if hasMin && hasMax { + query += " AND " + } + + if hasMax { + query += fmt.Sprintf(`"%s"`, sourceDbInfo.PrimaryKey) + if isMaxInclusive { + query += " <=" + } else { + query += " <" + } + query += fmt.Sprintf(" $%d", paramIdx) + } + } + + query += fmt.Sprintf(` ORDER BY "%s" ASC`, sourceDbInfo.PrimaryKey) + + return query } -func (postgresEx *PostgresExtractor) ProcessPartition( +func (postgresEx *PostgresExtractor) Exec( ctx context.Context, tableInfo config.SourceTableInfo, columns []models.ColumnType, @@ -60,14 +100,20 @@ func (postgresEx *PostgresExtractor) ProcessPartition( indexPrimaryKey int, chBatchesOut chan<- models.Batch, ) (int, error) { - query := buildExtractQueryPostgres(tableInfo, columns) + hasMin := partition.HasRange && partition.Range.Min > 0 + hasMax := partition.HasRange && partition.Range.Max > 0 + query := buildExtractQueryPostgres(tableInfo, columns, partition.HasRange, partition.Range.IsMinInclusive, partition.Range.IsMaxInclusive, hasMin, hasMax) - if partition.HasRange { - return 0, errors.New("Batch config not yet supported") + var queryArgs []any + if hasMin { + queryArgs = append(queryArgs, partition.Range.Min) + } + if hasMax { + queryArgs = append(queryArgs, partition.Range.Max) } rowsRead := 0 - rows, err := postgresEx.db.Query(ctx, query) + rows, err := postgresEx.db.Query(ctx, query, queryArgs...) if err != nil { return rowsRead, &custom_errors.ExtractorError{Partition: partition, HasLastId: false, Msg: err.Error()} } @@ -78,7 +124,7 @@ func (postgresEx *PostgresExtractor) ProcessPartition( for rows.Next() { values, err := rows.Values() if err != nil { - return rowsRead, errors.New("Unexpected error reading rows from source") + return rowsRead, &custom_errors.ExtractorError{Partition: partition, HasLastId: false, Msg: err.Error()} } rowsRead++ @@ -96,7 +142,7 @@ func (postgresEx *PostgresExtractor) ProcessPartition( } if err := rows.Err(); err != nil { - return rowsRead, errors.New("Unexpected error reading rows from source") + return rowsRead, &custom_errors.ExtractorError{Partition: partition, HasLastId: false, Msg: err.Error()} } if len(batchRows) > 0 { @@ -109,17 +155,3 @@ func (postgresEx *PostgresExtractor) ProcessPartition( return rowsRead, nil } - -func (postgresEx *PostgresExtractor) Exec( - ctx context.Context, - tableInfo config.SourceTableInfo, - columns []models.ColumnType, - batchSize int, - chPartitionsIn <-chan models.Partition, - chBatchesOut chan<- models.Batch, - chErrorsOut chan<- custom_errors.ExtractorError, - chJobErrorsOut chan<- custom_errors.JobError, - wgActivePartitions *sync.WaitGroup, - rowsRead *int64, -) { -} diff --git a/internal/app/etl/loaders/postgres.go b/internal/app/etl/loaders/main.go similarity index 81% rename from internal/app/etl/loaders/postgres.go rename to internal/app/etl/loaders/main.go index f9ed76f..523d6c8 100644 --- a/internal/app/etl/loaders/postgres.go +++ b/internal/app/etl/loaders/main.go @@ -15,31 +15,21 @@ import ( "github.com/jackc/pgx/v5/pgconn" ) -type PostgresLoader struct { +type GenericLoader struct { db dbwrapper.DbWrapper } -func NewPostgresLoader(db dbwrapper.DbWrapper) etl.Loader { - return &PostgresLoader{db: db} +func NewGenericLoader(db dbwrapper.DbWrapper) etl.Loader { + return &GenericLoader{db: db} } -func mapSlice[T any, V any](input []T, mapper func(T) V) []V { - result := make([]V, len(input)) - - for i, v := range input { - result[i] = mapper(v) - } - - return result -} - -func (postgresLd *PostgresLoader) ProcessBatch( +func (gl *GenericLoader) ProcessBatch( ctx context.Context, tableInfo config.TargetTableInfo, colNames []string, batch models.Batch, ) (int, error) { - _, err := postgresLd.db.SaveMassive( + _, err := gl.db.SaveMassive( ctx, tableInfo.Schema, tableInfo.Table, @@ -65,7 +55,7 @@ func (postgresLd *PostgresLoader) ProcessBatch( return len(batch.Rows), nil } -func (postgresLd *PostgresLoader) Exec( +func (gl *GenericLoader) Exec( ctx context.Context, tableInfo config.TargetTableInfo, columns []models.ColumnType, @@ -92,7 +82,7 @@ func (postgresLd *PostgresLoader) Exec( return } - processedRows, err := postgresLd.ProcessBatch(ctx, tableInfo, colNames, batch) + processedRows, err := gl.ProcessBatch(ctx, tableInfo, colNames, batch) if err != nil { var ldError *custom_errors.LoaderError diff --git a/internal/app/etl/loaders/types.go b/internal/app/etl/loaders/types.go deleted file mode 100644 index c88d5fb..0000000 --- a/internal/app/etl/loaders/types.go +++ /dev/null @@ -1 +0,0 @@ -package loaders diff --git a/internal/app/etl/loaders/utils.go b/internal/app/etl/loaders/utils.go new file mode 100644 index 0000000..1a2dbf4 --- /dev/null +++ b/internal/app/etl/loaders/utils.go @@ -0,0 +1,11 @@ +package loaders + +func mapSlice[T any, V any](input []T, mapper func(T) V) []V { + result := make([]V, len(input)) + + for i, v := range input { + result[i] = mapper(v) + } + + return result +} diff --git a/internal/app/etl/table_analyzers/main.go b/internal/app/etl/table_analyzers/main.go index c226dbe..cb08390 100644 --- a/internal/app/etl/table_analyzers/main.go +++ b/internal/app/etl/table_analyzers/main.go @@ -15,7 +15,22 @@ func PartitionRangeGenerator( tableInfo config.TableInfo, partitionColumn string, rowsPerPartition int64, + jobRange config.RangeConfig, ) ([]models.Partition, error) { + if jobRange.Min > 0 { + return []models.Partition{{ + Id: uuid.New(), + HasRange: true, + RetryCounter: 0, + Range: models.PartitionRange{ + Min: jobRange.Min, + Max: jobRange.Max, + IsMinInclusive: jobRange.IsMinInclusive, + IsMaxInclusive: jobRange.IsMaxInclusive, + }, + }}, nil + } + rowsCount, err := tableAnalyzer.EstimateTotalRows(ctx, tableInfo) if err != nil { return nil, err diff --git a/internal/app/etl/table_analyzers/mssql.go b/internal/app/etl/table_analyzers/mssql.go index 0ec0f12..4faaf21 100644 --- a/internal/app/etl/table_analyzers/mssql.go +++ b/internal/app/etl/table_analyzers/mssql.go @@ -36,7 +36,7 @@ JOIN sys.types t ON c.user_type_id = t.user_type_id LEFT JOIN sys.types bt ON t.is_user_defined = 1 AND bt.user_type_id = t.system_type_id JOIN sys.tables st ON c.object_id = st.object_id JOIN sys.schemas s ON st.schema_id = s.schema_id -WHERE s.name = @schema AND st.name = @table AND c.name NOT LIKE 'graph_id%' +WHERE s.name = @schema AND st.name = @table AND (c.is_hidden = 0 OR (c.graph_type IS NOT NULL AND c.name LIKE '$%')) ORDER BY c.column_id;` type rawColumnMssql struct { @@ -184,7 +184,7 @@ JOIN sys.partitions p ON t.object_id = p.object_id WHERE s.name = @schema AND t.name = @table AND p.index_id IN (0, 1) GROUP BY t.name` - ctxTimeout, cancel := context.WithTimeout(ctx, time.Second*20) + ctxTimeout, cancel := context.WithTimeout(ctx, 1*time.Minute) defer cancel() var rowsCount int64 @@ -216,7 +216,7 @@ ORDER BY batch_id`, tableInfo.Schema, tableInfo.Table) - ctxTimeout, cancel := context.WithTimeout(ctx, time.Second*20) + ctxTimeout, cancel := context.WithTimeout(ctx, 1*time.Minute) defer cancel() rows, err := ta.db.Query(ctxTimeout, query, sql.Named("maxPartitions", maxPartitions)) diff --git a/internal/app/etl/table_analyzers/postgres.go b/internal/app/etl/table_analyzers/postgres.go index 9489c78..8aac15d 100644 --- a/internal/app/etl/table_analyzers/postgres.go +++ b/internal/app/etl/table_analyzers/postgres.go @@ -125,7 +125,7 @@ func (ta *PostgresTableAnalyzer) QueryColumnTypes( ctx context.Context, tableInfo config.TableInfo, ) ([]models.ColumnType, error) { - localCtx, cancel := context.WithTimeout(ctx, 20*time.Second) + localCtx, cancel := context.WithTimeout(ctx, 1*time.Minute) defer cancel() rows, err := ta.db.Query(localCtx, postgresColumnMetadataQuery, tableInfo.Schema, tableInfo.Table) diff --git a/internal/app/etl/transformers/mssql.go b/internal/app/etl/transformers/mssql.go index 7270ebb..0500a6e 100644 --- a/internal/app/etl/transformers/mssql.go +++ b/internal/app/etl/transformers/mssql.go @@ -3,18 +3,32 @@ package transformers import ( "context" "errors" + "fmt" + "strings" "sync" "time" + "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/azure" + "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/config" "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/custom_errors" "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/etl" "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/models" + "github.com/google/uuid" + log "github.com/sirupsen/logrus" ) -type MssqlTransformer struct{} +type MssqlTransformer struct { + toStorage config.ToStorageConfig + sourceTable config.SourceTableInfo + azureClient *azure.Client +} -func NewMssqlTransformer() etl.Transformer { - return &MssqlTransformer{} +func NewMssqlTransformer(toStorage config.ToStorageConfig, sourceTable config.SourceTableInfo, azureClient *azure.Client) etl.Transformer { + return &MssqlTransformer{ + toStorage: toStorage, + sourceTable: sourceTable, + azureClient: azureClient, + } } func computeTransformationPlan(columns []models.ColumnType) []etl.ColumnTransformPlan { @@ -60,6 +74,63 @@ func computeTransformationPlan(columns []models.ColumnType) []etl.ColumnTransfor return plan } +func computeStorageTransformationPlan( + ctx context.Context, + azureClient *azure.Client, + toStorage config.ToStorageConfig, + sourceColumns []models.ColumnType, + sourceTable config.SourceTableInfo, +) []etl.ColumnTransformPlan { + if azureClient == nil || len(toStorage.Columns) == 0 { + return nil + } + + colIndex := make(map[string]int, len(sourceColumns)) + for i, col := range sourceColumns { + colIndex[strings.ToUpper(col.Name())] = i + } + + var plan []etl.ColumnTransformPlan + for _, storageCol := range toStorage.Columns { + if storageCol.Mode != "REFERENCE_ONLY" { + log.Warnf("to_storage: unsupported mode %q for column %s — skipping", storageCol.Mode, storageCol.Source) + continue + } + + idx, ok := colIndex[strings.ToUpper(storageCol.Source)] + if !ok { + log.Warnf("to_storage: source column %q not found in source schema — skipping", storageCol.Source) + continue + } + + sourceColName := storageCol.Source + schema := sourceTable.Schema + table := sourceTable.Table + + plan = append(plan, etl.ColumnTransformPlan{ + Index: idx, + Fn: func(v any) (any, error) { + if v == nil { + return nil, nil + } + b, ok := v.([]byte) + if !ok { + log.Warnf("to_storage: expected []byte for %s.%s.%s, got %T — passing through", + schema, table, sourceColName, v) + return v, nil + } + blobPath := fmt.Sprintf("%s/%s/%s/%s.bin", schema, table, sourceColName, uuid.New().String()) + blobURL, err := azureClient.UploadAndGetURL(ctx, blobPath, b) + if err != nil { + return nil, fmt.Errorf("uploading %s.%s.%s: %w", schema, table, sourceColName, err) + } + return blobURL, nil + }, + }) + } + return plan +} + const processBatchCtxCheck = 4096 func (mssqlTr *MssqlTransformer) ProcessBatch( @@ -100,6 +171,8 @@ func (mssqlTr *MssqlTransformer) Exec( wgActiveBatches *sync.WaitGroup, ) { transformationPlan := computeTransformationPlan(columns) + storagePlan := computeStorageTransformationPlan(ctx, mssqlTr.azureClient, mssqlTr.toStorage, columns, mssqlTr.sourceTable) + transformationPlan = append(transformationPlan, storagePlan...) for { if ctx.Err() != nil { diff --git a/internal/app/etl/types.go b/internal/app/etl/types.go index e6c8ee5..f56f8d1 100644 --- a/internal/app/etl/types.go +++ b/internal/app/etl/types.go @@ -10,7 +10,7 @@ import ( ) type Extractor interface { - ProcessPartition( + Exec( ctx context.Context, tableInfo config.SourceTableInfo, columns []models.ColumnType, @@ -19,19 +19,6 @@ type Extractor interface { indexPrimaryKey int, chBatchesOut chan<- models.Batch, ) (int, error) - - Exec( - ctx context.Context, - tableInfo config.SourceTableInfo, - columns []models.ColumnType, - batchSize int, - chPartitionsIn <-chan models.Partition, - chBatchesOut chan<- models.Batch, - chErrorsOut chan<- custom_errors.ExtractorError, - chJobErrorsOut chan<- custom_errors.JobError, - wgActivePartitions *sync.WaitGroup, - rowsRead *int64, - ) } type TransformerFunc func(any) (any, error) diff --git a/internal/app/models/main.go b/internal/app/models/main.go index 6156a86..60eb73e 100644 --- a/internal/app/models/main.go +++ b/internal/app/models/main.go @@ -1,6 +1,10 @@ package models -import "github.com/google/uuid" +import ( + "time" + + "github.com/google/uuid" +) type UnknownRowValues = []any @@ -25,3 +29,13 @@ type Partition struct { HasRange bool RetryCounter int } + +type JobResult struct { + JobName string + StartTime time.Time + Duration time.Duration + RowsRead int64 + RowsLoaded int64 + RowsFailed int64 + Error error +}