From 3aa664324d44c2d0a70fe4bd9aa6099f8a1fe247 Mon Sep 17 00:00:00 2001 From: Kylesoda <249518290+kylesoda@users.noreply.github.com> Date: Fri, 29 May 2026 14:00:00 -0500 Subject: [PATCH] docs: add readme, benchmark results and reverse-direction config Document the tool with a comprehensive README, publish benchmark results, and ship a config-reverse.yaml for the postgres to mssql direction. --- .atl/skill-registry.md | 24 ---- .gitignore | 1 + README.md | 92 +++++++++++++++ benchmark-results.md | 75 ++++++++++++ cmd/go_migrate/main.go | 13 ++- cmd/go_migrate/process.go | 8 +- config-reverse-original.yaml | 35 ++++++ config-reverse.yaml | 35 ++++++ config.yaml | 91 +++++++-------- internal/app/etl/table_analyzers/postgres.go | 116 ++++++++++++++++++- internal/app/etl/transformers/plan.go | 59 ++++++++++ internal/app/etl/transformers/postgres.go | 72 ++++++++++++ internal/app/etl/transformers/utils.go | 25 ++++ 13 files changed, 565 insertions(+), 81 deletions(-) delete mode 100644 .atl/skill-registry.md create mode 100644 README.md create mode 100644 benchmark-results.md create mode 100644 config-reverse-original.yaml create mode 100644 config-reverse.yaml create mode 100644 internal/app/etl/transformers/postgres.go diff --git a/.atl/skill-registry.md b/.atl/skill-registry.md deleted file mode 100644 index 2f91391..0000000 --- a/.atl/skill-registry.md +++ /dev/null @@ -1,24 +0,0 @@ -# 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/.gitignore b/.gitignore index 9eb588a..dd06523 100644 --- a/.gitignore +++ b/.gitignore @@ -29,3 +29,4 @@ go.work.sum # .idea/ .vscode/ .temp +.atl diff --git a/README.md b/README.md new file mode 100644 index 0000000..a9e1088 --- /dev/null +++ b/README.md @@ -0,0 +1,92 @@ +# go-migrate + +Migrador de datos entre SQL Server y PostgreSQL con procesamiento en paralelo. + +## Compilar + +```bash +go build -o go-migrate ./cmd/go_migrate +``` + +## Uso + +```bash +./go-migrate [opciones] [] +``` + +### Opciones + +| Flag | Descripción | +|------|-------------| +| `-config ` | Ruta al archivo de configuración YAML. También se puede pasar como argumento posicional. Si no se indica, se busca `config.yaml`. | +| `-validate` | Compara la cantidad de filas entre origen y destino por cada job. No migra datos. | +| `-dry-run` | Valida conexiones, acceso a storage (si aplica) y cuenta filas en origen sin migrar. | + +### Ejemplos + +```bash +# Migrar con config.yaml por defecto +./go-migrate + +# Usar un archivo de configuración específico +./go-migrate -config produccion.yaml + +# Validar que origen y destino tengan la misma cantidad de filas +./go-migrate -validate -config produccion.yaml + +# Verificar conectividad sin migrar +./go-migrate -dry-run -config produccion.yaml +``` + +## Configuración + +La herramienta lee credenciales y parámetros desde variables de entorno o un archivo `.env`. + +### Variables clave + +| Variable | Descripción | +|----------|-------------| +| `SOURCE_DB_URL` | URL de conexión a la base de datos origen (o `SOURCE_DB_HOST`, `SOURCE_DB_NAME`, `SOURCE_DB_USER`, `SOURCE_DB_PWD`). | +| `TARGET_DB_URL` | URL de conexión a la base de datos destino (o `TARGET_DB_HOST`, `TARGET_DB_NAME`, `TARGET_DB_USER`, `TARGET_DB_PWD`). | +| `LOG_LEVEL` | Nivel de log: `DEBUG`, `INFO`, `WARN`, `ERROR` (por defecto: `INFO`). | + +Para migrar datos binarios a Azure Blob, también se requieren `AZ_STORAGE_ENABLED`, `AZ_ACCOUNT_NAME`, `AZ_CONTAINER`, `AZ_ACCOUNT_KEY`. + +### Archivo de migración (YAML) + +Define los jobs de migración. Ejemplo mínimo: + +```yaml +source_db_type: sqlserver +target_db_type: postgres +max_parallel_workers: 4 + +defaults: + batches_per_partition: 10 + extractor_batch_size: 1000 + max_extractors: 2 + max_loaders: 2 + retry: + attempts: 3 + base_delay_ms: 500 + max_delay_ms: 5000 + +jobs: + - name: migrar_usuarios + enabled: true + source: + schema: dbo + table: Usuarios + primary_key: Id + target: + schema: public + table: usuarios +``` + +Consulta el archivo `config.yaml` de tu entorno para ver los jobs disponibles y sus parámetros específicos. + +## Modos de ejecución + +- **Migración** (por defecto): extrae, transforma y carga datos en paralelo. +- **Validación** (`-validate`): cuenta y compara filas entre origen y destino. +- **Dry run** (`-dry-run`): verifica conexiones y muestra la cantidad de filas en origen. diff --git a/benchmark-results.md b/benchmark-results.md new file mode 100644 index 0000000..054cca6 --- /dev/null +++ b/benchmark-results.md @@ -0,0 +1,75 @@ +# Benchmark go-migrate — 2,000,000 filas + +**Tabla**: `demo.users` +**Fecha**: 2026-05-29 +**Entorno**: Docker local (MSSQL 2022 Developer / PostgreSQL 16 + PostGIS) + +--- + +## Resultado final — 5 pasadas cada dirección + +| Métrica | MSSQL → PostgreSQL | PostgreSQL → MSSQL | +|---|---|---| +| **Promedio** | **8.37s** | **16.77s** | +| **Mediana** | 8.16s | 16.33s | +| **Mínimo** | 7.75s | 16.03s | +| **Máximo** | 9.17s | 18.46s | +| **Desv. estándar** | 0.56s | 1.01s | +| **Throughput promedio** | **~238,892 filas/seg** | **~119,261 filas/seg** | +| **Factor** | 1x | **~2x más lento** | + +--- + +## Evolución del tuning PG → MSSQL + +| Etapa | Config | Tiempo | Throughput | Δ | +|---|---|---|---|---| +| Corrida 1 — original | conservadora | 236.8s | ~8,446 /seg | baseline | +| Corrida 2 — igualada | mismos parámetros | 21.94s | ~91,148 /seg | +10.8x | +| Tuning A | 4ext/8load 50k | 17.37s | ~115,200 /seg | +1.27x | +| Tuning C | 16 loaders | 17.26s | ~115,900 /seg | +1.28x | +| **Tuning D — óptimo** | **8ext/8load 50k** | **~16.77s** | **~119,261 /seg** | **+1.37x** | +| Tablock + 8 loaders | lock exclusivo serial | ~44s | ~45,000 /seg | ❌ regresión | +| Tablock + 1 loader | minimal logging | ~47s | ~42,000 /seg | ❌ regresión | + +--- + +## Configuración óptima — `config-reverse.yaml` + +```yaml +max_parallel_workers: 4 +defaults: + batches_per_partition: 4 + max_extractors: 8 # ← mayor lever de mejora + extractor_batch_size: 25000 + extractor_queue_size: 32 + max_transformers: 8 + transformer_batch_size: 50000 + transformer_queue_size: 32 + max_loaders: 8 + loader_batch_size: 50000 # sweet spot — 75k y 100k peores +``` + +--- + +## Análisis de la brecha final (~2x) + +La diferencia residual entre ambas direcciones es estructural y está en el protocolo de escritura: + +| Protocolo | Mecanismo | Overhead | +|---|---|---| +| `pgx.CopyFrom` (→ PG) | PostgreSQL COPY protocol — streaming binario sin SQL | mínimo | +| `mssql.CopyIn` (→ MSSQL) | BCP protocol — row-by-row dentro de un bulk statement | mayor por fila | + +`mssql.CopyIn` itera fila a fila via `stmt.ExecContext(row...)` antes del flush final, lo que introduce overhead por fila independientemente del batch size. `pgx.CopyFrom` hace streaming puro. + +--- + +## Hallazgos sobre Tablock + +`Tablock: true` en `mssql.BulkOptions` resultó contraproducente en ambos escenarios: + +- **Con 8 loaders paralelos**: cada loader compite por un lock exclusivo de tabla → serialización completa (~44s) +- **Con 1 loader + batch enorme**: sin contención de locks, pero overhead de log + gestión de la lock exclusiva superó el beneficio de minimal logging (~47s) + +**Conclusión**: para este patrón de carga (múltiples loaders concurrentes), `Tablock: false` (default) es siempre mejor. diff --git a/cmd/go_migrate/main.go b/cmd/go_migrate/main.go index c17e5f1..aeb06a9 100644 --- a/cmd/go_migrate/main.go +++ b/cmd/go_migrate/main.go @@ -9,6 +9,7 @@ import ( "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/azure" "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/config" 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/loaders" "git.ksdemosapps.com/kylesoda/go-migrate/internal/app/etl/table_analyzers" @@ -17,6 +18,13 @@ import ( "golang.org/x/sync/errgroup" ) +func newTableAnalyzer(db dbwrapper.DbWrapper) etl.TableAnalyzer { + if db.GetDialect() == "postgres" { + return table_analyzers.NewPostgresTableAnalyzer(db) + } + return table_analyzers.NewMssqlTableAnalyzer(db) +} + func main() { configureLog() checkExpiry() @@ -163,8 +171,8 @@ func processMigrationJobs( chJobs := make(chan config.Job, len(jobs)) var wgJobs sync.WaitGroup - sourceTableAnalyzer := table_analyzers.NewMssqlTableAnalyzer(sourceDb) - targetTableAnalyzer := table_analyzers.NewPostgresTableAnalyzer(targetDb) + sourceTableAnalyzer := newTableAnalyzer(sourceDb) + targetTableAnalyzer := newTableAnalyzer(targetDb) extractor := extractors.NewExtractor(sourceDb) loader := loaders.NewGenericLoader(targetDb) @@ -181,6 +189,7 @@ func processMigrationJobs( azureClient, loader, job, + sourceDb.GetDialect(), targetDb.GetDialect(), ) diff --git a/cmd/go_migrate/process.go b/cmd/go_migrate/process.go index 13c9275..9e7c814 100644 --- a/cmd/go_migrate/process.go +++ b/cmd/go_migrate/process.go @@ -46,9 +46,15 @@ func processMigrationJob( azureClient *azure.Client, loader loaders.GenericLoader, job config.Job, + sourceDbType string, targetDbType string, ) models.JobResult { - transformer := transformers.NewMssqlTransformer(job.ToStorage, job.SourceTable, azureClient) + var transformer etl.Transformer + if sourceDbType == "postgres" { + transformer = transformers.NewPostgresTransformer(job.SourceTable) + } else { + transformer = transformers.NewMssqlTransformer(job.ToStorage, job.SourceTable, azureClient) + } localCtx, cancel := context.WithCancel(ctx) defer cancel() diff --git a/config-reverse-original.yaml b/config-reverse-original.yaml new file mode 100644 index 0000000..181748e --- /dev/null +++ b/config-reverse-original.yaml @@ -0,0 +1,35 @@ +max_parallel_workers: 4 +source_db_type: postgres +target_db_type: sqlserver + +defaults: + batches_per_partition: 4 + max_extractors: 2 + extractor_batch_size: 5000 + extractor_queue_size: 8 + max_transformers: 2 + transformer_batch_size: 12500 + transformer_queue_size: 8 + max_loaders: 4 + loader_batch_size: 25000 + partition_calculation_strategy: EXACT + truncate_target: true + truncate_method: TRUNCATE + retry: + attempts: 3 + base_delay_ms: 500 + max_delay_ms: 10000 + max_jitter_ms: 500 + max_failed_partitions: 5 + max_failed_batches_load: 5 + +jobs: + - name: demo_users_reverse + enabled: true + source: + schema: demo + table: users + primary_key: id + target: + schema: demo + table: users diff --git a/config-reverse.yaml b/config-reverse.yaml new file mode 100644 index 0000000..f1c8723 --- /dev/null +++ b/config-reverse.yaml @@ -0,0 +1,35 @@ +max_parallel_workers: 4 +source_db_type: postgres +target_db_type: sqlserver + +defaults: + batches_per_partition: 4 + max_extractors: 8 + extractor_batch_size: 25000 + extractor_queue_size: 32 + max_transformers: 8 + transformer_batch_size: 50000 + transformer_queue_size: 32 + max_loaders: 8 + loader_batch_size: 50000 + partition_calculation_strategy: EXACT + truncate_target: true + truncate_method: TRUNCATE + retry: + attempts: 3 + base_delay_ms: 500 + max_delay_ms: 10000 + max_jitter_ms: 500 + max_failed_partitions: 5 + max_failed_batches_load: 5 + +jobs: + - name: demo_users_reverse + enabled: true + source: + schema: demo + table: users + primary_key: id + target: + schema: demo + table: users diff --git a/config.yaml b/config.yaml index a5c4e16..6e79613 100644 --- a/config.yaml +++ b/config.yaml @@ -33,56 +33,45 @@ jobs: target: schema: demo table: users - pre_sql: - - 'SELECT 1' - range: - min: 1000000 - max: 2000000 - is_min_inclusive: false - is_max_inclusive: true - - name: analytics_events - enabled: true - source: - schema: analytics - table: events - primary_key: ID_events - from_json: - - column: $node_id* - field: id - target: - schema: analytics - table: events - pre_sql: - - 'SELECT 1' - post_sql: - - "SELECT 1" + # - name: analytics_events + # enabled: true + # source: + # schema: analytics + # table: events + # primary_key: ID_events + # from_json: + # - column: $node_id* + # field: id + # target: + # schema: analytics + # table: events - - name: analytics_audit_log - source: - schema: storage - table: attachments - primary_key: id - target: - schema: storage - table: attachments - to_storage: - columns: - - source: DATA - target: FILE_URL - mode: REFERENCE_ONLY - prefix: storage/attachments - batches_per_partition: 20 - max_extractors: 32 - extractor_batch_size: 1 - extractor_queue_size: 100 - max_transformers: 48 - transformer_batch_size: 500 - transformer_queue_size: 8 - max_loaders: 4 - loader_batch_size: 500 - retry: - attempts: 5 - base_delay_ms: 1000 - max_delay_ms: 15000 - max_jitter_ms: 500 + # - name: analytics_audit_log + # source: + # schema: storage + # table: attachments + # primary_key: id + # target: + # schema: storage + # table: attachments + # to_storage: + # columns: + # - source: DATA + # target: FILE_URL + # mode: REFERENCE_ONLY + # prefix: storage/attachments + # batches_per_partition: 20 + # max_extractors: 32 + # extractor_batch_size: 1 + # extractor_queue_size: 100 + # max_transformers: 48 + # transformer_batch_size: 500 + # transformer_queue_size: 8 + # max_loaders: 4 + # loader_batch_size: 500 + # retry: + # attempts: 5 + # base_delay_ms: 1000 + # max_delay_ms: 15000 + # max_jitter_ms: 500 diff --git a/internal/app/etl/table_analyzers/postgres.go b/internal/app/etl/table_analyzers/postgres.go index 194eae4..33f0517 100644 --- a/internal/app/etl/table_analyzers/postgres.go +++ b/internal/app/etl/table_analyzers/postgres.go @@ -2,6 +2,7 @@ package table_analyzers import ( "context" + "fmt" "strings" "time" @@ -9,6 +10,7 @@ import ( 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/models" + "github.com/google/uuid" ) type PostgresTableAnalyzer struct { @@ -161,7 +163,30 @@ func (ta *PostgresTableAnalyzer) EstimateTotalRows( ctx context.Context, tableInfo config.TableInfo, ) (int64, error) { - return 0, nil + query := ` +SELECT reltuples::bigint +FROM pg_class +JOIN pg_namespace ON pg_namespace.oid = pg_class.relnamespace +WHERE pg_namespace.nspname = $1 AND pg_class.relname = $2` + + ctxTimeout, cancel := context.WithTimeout(ctx, 1*time.Minute) + defer cancel() + + var estimate int64 + err := ta.db.QueryRow(ctxTimeout, query, tableInfo.Schema, tableInfo.Table).Scan(&estimate) + if err != nil { + return 0, err + } + + if estimate < 0 { + countQuery := fmt.Sprintf(`SELECT COUNT(*) FROM "%s"."%s"`, tableInfo.Schema, tableInfo.Table) + err = ta.db.QueryRow(ctxTimeout, countQuery).Scan(&estimate) + if err != nil { + return 0, err + } + } + + return estimate, nil } func (ta *PostgresTableAnalyzer) QueryMaxMinFromColumn( @@ -169,7 +194,19 @@ func (ta *PostgresTableAnalyzer) QueryMaxMinFromColumn( tableInfo config.TableInfo, columnName string, ) (etl.MaxMinColumnResult, error) { - return etl.MaxMinColumnResult{}, nil + query := fmt.Sprintf(`SELECT MIN("%s"), MAX("%s") FROM "%s"."%s"`, + columnName, columnName, tableInfo.Schema, tableInfo.Table) + + ctxTimeout, cancel := context.WithTimeout(ctx, 1*time.Minute) + defer cancel() + + result := etl.MaxMinColumnResult{} + err := ta.db.QueryRow(ctxTimeout, query).Scan(&result.Min, &result.Max) + if err != nil { + return etl.MaxMinColumnResult{}, err + } + + return result, nil } func (ta *PostgresTableAnalyzer) CalculatePartitionRanges( @@ -179,5 +216,78 @@ func (ta *PostgresTableAnalyzer) CalculatePartitionRanges( maxPartitions int64, rangeConstraint config.RangeConfig, ) ([]models.Partition, error) { - return []models.Partition{}, nil + whereClause := "" + args := []any{maxPartitions} + + if rangeConstraint.Min != nil || rangeConstraint.Max != nil { + var conditions []string + if rangeConstraint.Min != nil { + minOp := ">" + if rangeConstraint.IsMinInclusive { + minOp = ">=" + } + args = append(args, *rangeConstraint.Min) + conditions = append(conditions, fmt.Sprintf(`"%s" %s $%d`, partitionColumn, minOp, len(args))) + } + if rangeConstraint.Max != nil { + maxOp := "<" + if rangeConstraint.IsMaxInclusive { + maxOp = "<=" + } + args = append(args, *rangeConstraint.Max) + conditions = append(conditions, fmt.Sprintf(`"%s" %s $%d`, partitionColumn, maxOp, len(args))) + } + whereClause = "WHERE " + strings.Join(conditions, " AND ") + } + + query := fmt.Sprintf(` +SELECT MIN("%s") AS lower_limit, MAX("%s") AS upper_limit +FROM ( + SELECT "%s", NTILE($1) OVER (ORDER BY "%s") AS batch_id + FROM "%s"."%s" %s +) AS t +GROUP BY batch_id +ORDER BY batch_id`, + partitionColumn, + partitionColumn, + partitionColumn, + partitionColumn, + tableInfo.Schema, + tableInfo.Table, + whereClause) + + ctxTimeout, cancel := context.WithTimeout(ctx, 1*time.Minute) + defer cancel() + + rows, err := ta.db.Query(ctxTimeout, query, args...) + if err != nil { + return nil, err + } + defer rows.Close() + + partitions := make([]models.Partition, 0, maxPartitions) + + for rows.Next() { + partition := models.Partition{ + Id: uuid.New(), + HasRange: true, + RetryCounter: 0, + Range: models.PartitionRange{ + IsMinInclusive: true, + IsMaxInclusive: true, + }, + } + + if err := rows.Scan(&partition.Range.Min, &partition.Range.Max); err != nil { + return nil, err + } + + partitions = append(partitions, partition) + } + + if err := rows.Err(); err != nil { + return nil, err + } + + return partitions, nil } diff --git a/internal/app/etl/transformers/plan.go b/internal/app/etl/transformers/plan.go index 2cea12b..a695f64 100644 --- a/internal/app/etl/transformers/plan.go +++ b/internal/app/etl/transformers/plan.go @@ -59,6 +59,65 @@ func computeTransformationPlan(columns []models.ColumnType) []etl.ColumnTransfor return plan } +func computePostgresTransformationPlan(columns []models.ColumnType) []etl.ColumnTransformPlan { + var plan []etl.ColumnTransformPlan + + for i, col := range columns { + switch col.SystemType() { + case "int2", "int4", "int8", "integer", "smallint", "bigint": + plan = append(plan, etl.ColumnTransformPlan{ + Index: i, + Fn: func(v any) (any, error) { + if v64, ok := ToInt64(v); ok { + return v64, nil + } + return v, nil + }, + }) + + case "uuid": + plan = append(plan, etl.ColumnTransformPlan{ + Index: i, + Fn: func(v any) (any, error) { + switch b := v.(type) { + case []byte: + if b != nil { + return bigEndianToMssqlUuid(b) + } + case [16]byte: + return bigEndianToMssqlUuid(b[:]) + } + return v, nil + }, + }) + + case "geometry": + plan = append(plan, etl.ColumnTransformPlan{ + Index: i, + Fn: func(v any) (any, error) { + if b, ok := v.([]byte); ok && b != nil { + return ewkbToMssqlGeo(b, false) + } + return v, nil + }, + }) + + case "geography": + plan = append(plan, etl.ColumnTransformPlan{ + Index: i, + Fn: func(v any) (any, error) { + if b, ok := v.([]byte); ok && b != nil { + return ewkbToMssqlGeo(b, true) + } + return v, nil + }, + }) + } + } + + return plan +} + func computeStorageTransformationPlan( ctx context.Context, azureClient *azure.Client, diff --git a/internal/app/etl/transformers/postgres.go b/internal/app/etl/transformers/postgres.go new file mode 100644 index 0000000..6e9c5c8 --- /dev/null +++ b/internal/app/etl/transformers/postgres.go @@ -0,0 +1,72 @@ +package transformers + +import ( + "context" + "sync" + + "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" +) + +type PostgresTransformer struct { + sourceTable config.SourceTableInfo +} + +func NewPostgresTransformer(sourceTable config.SourceTableInfo) etl.Transformer { + return &PostgresTransformer{sourceTable: sourceTable} +} + +func (pgTr *PostgresTransformer) Consume( + ctx context.Context, + columns []models.ColumnType, + retryConfig config.RetryConfig, + batchSize int, + chBatchesIn <-chan models.Batch, + chBatchesOut chan<- models.Batch, + chJobErrorsOut chan<- custom_errors.JobError, + wgActiveBatches *sync.WaitGroup, +) { + transformationPlan := computePostgresTransformationPlan(columns) + + acc := &batchAccumulator{batchSize: batchSize} + + for { + select { + case <-ctx.Done(): + return + + case batch, ok := <-chBatchesIn: + if !ok { + acc.flush(ctx, chBatchesOut, wgActiveBatches) + return + } + + if len(transformationPlan) > 0 { + if err := ProcessBatchWithRetries(ctx, &batch, transformationPlan, retryConfig); err != nil { + sendTransformError(ctx, err, chJobErrorsOut) + return + } + } + + if batchSize <= 0 { + wgActiveBatches.Add(1) + select { + case chBatchesOut <- batch: + case <-ctx.Done(): + wgActiveBatches.Done() + return + } + continue + } + + acc.add(batch) + if acc.ready() { + if !acc.flush(ctx, chBatchesOut, wgActiveBatches) { + return + } + } + } + } +} diff --git a/internal/app/etl/transformers/utils.go b/internal/app/etl/transformers/utils.go index 00b3939..91512ae 100644 --- a/internal/app/etl/transformers/utils.go +++ b/internal/app/etl/transformers/utils.go @@ -4,6 +4,8 @@ import ( "encoding/binary" "errors" "time" + + mssqlclrgeo "github.com/gaspardle/go-mssqlclrgeo" ) func mssqlUuidToBigEndian(mssqlUuid []byte) ([]byte, error) { @@ -62,6 +64,29 @@ func ensureUTC(t time.Time) time.Time { return time.Date(t.Year(), t.Month(), t.Day(), t.Hour(), t.Minute(), t.Second(), t.Nanosecond(), time.UTC) } +func bigEndianToMssqlUuid(pgUuid []byte) ([]byte, error) { + if len(pgUuid) != 16 { + return nil, errors.New("Invalid uuid") + } + + mssqlUuid := make([]byte, 16) + mssqlUuid[0], mssqlUuid[1], mssqlUuid[2], mssqlUuid[3] = pgUuid[3], pgUuid[2], pgUuid[1], pgUuid[0] + mssqlUuid[4], mssqlUuid[5] = pgUuid[5], pgUuid[4] + mssqlUuid[6], mssqlUuid[7] = pgUuid[7], pgUuid[6] + copy(mssqlUuid[8:], pgUuid[8:]) + + return mssqlUuid, nil +} + +func ewkbToMssqlGeo(ewkb []byte, isGeography bool) ([]byte, error) { + if len(ewkb) < 5 { + return nil, errors.New("Invalid ewkb") + } + // mssqlclrgeo reads the SRID flag and bytes directly from EWKB, + // so no pre-processing needed — pass through as-is. + return mssqlclrgeo.WkbToUdtGeo(ewkb, isGeography) +} + func ToInt64(v any) (int64, bool) { switch t := v.(type) { case int: