/
iezhelev
/
cppsh_micro
Обзор
Документация
Войти
/
iezhelev
/
cppsh_micro
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
Безопасность
main
internal/sync/table_sync.go
221 строка
7 KB
iezhelev
updatedtime as is
03 июн 2025, 07:39
03 июн 2025, 07:39
9b17d7a
Код
Авторство
О чём код?
package sync import ( "context" "database/sql" "fmt" "log" "strings" "time" bd "cppsh_micro/internal/database/sync" ) const ( maxPostgresParams = 65535 // PostgreSQL's maximum parameters per query targetBatchParams = 60000 // Safe margin below the limit ) func SyncDataGeneric( ctx context.Context, mainDB, dockerDB *sql.DB, lastSyncTime time.Time, srcTableSchema string, srcTableName string, dstTableName string, additionalWhereClause string, ) (time.Time, error) { now := time.Now().UTC() srcFields := bd.CdcTableFieldsMapping[srcTableName].SrcFields dstFields := bd.CdcTableFieldsMapping[srcTableName].DstFields // Calculate dynamic batch size based on column count columns := len(dstFields) dynamicBatchSize := targetBatchParams / columns query := fmt.Sprintf( `SELECT %s FROM %s.%s WHERE COALESCE(updated_time, created_time) > $1 AT TIME ZONE 'UTC' AND COALESCE(updated_time, created_time) <= $2 AT TIME ZONE 'UTC'`, strings.Join(srcFields, ", "), srcTableSchema, srcTableName, ) if additionalWhereClause != "" { query += " " + additionalWhereClause } rows, err := mainDB.QueryContext(ctx, query, lastSyncTime, now) if err != nil { return lastSyncTime, fmt.Errorf("error fetching data: %w query %s", err, query) } defer rows.Close() tx, err := dockerDB.BeginTx(ctx, nil) if err != nil { return lastSyncTime, fmt.Errorf("failed to begin transaction: %w", err) } defer tx.Rollback() // Process rows in batches var currentBatch [][]interface{} var totalRecords int scanArgs := make([]interface{}, len(srcFields)) record := make([]interface{}, len(srcFields)) for i := range scanArgs { scanArgs[i] = &record[i] } for rows.Next() { err := rows.Scan(scanArgs...) if err != nil { log.Printf("Skipping row due to scan error: %v", err) continue } values := make([]interface{}, len(record)) copy(values, record) currentBatch = append(currentBatch, values) if len(currentBatch) >= dynamicBatchSize { if err := processBatch(ctx, tx, srcTableName, dstTableName, currentBatch); err != nil { return lastSyncTime, err } totalRecords += len(currentBatch) //log.Printf("batched save for table %s : %d", srcTableName, totalRecords) currentBatch = nil } } if len(currentBatch) > 0 { if err := processBatch(ctx, tx, srcTableName, dstTableName, currentBatch); err != nil { return lastSyncTime, err } totalRecords += len(currentBatch) } if err := tx.Commit(); err != nil { return lastSyncTime, fmt.Errorf("commit failed: %w", err) } if totalRecords > 0 { log.Printf("Synced %d records to %s", totalRecords, dstTableName) } return now, nil } func deduplicateBatch(batch [][]interface{}, srcTableName string) [][]interface{} { pkIndex := -1 // Find the index of the primary key in dstFields for i, field := range bd.CdcTableFieldsMapping[srcTableName].DstFields { if field == bd.CdcTableFieldsMapping[srcTableName].PrimaryKey { pkIndex = i break } } if pkIndex == -1 { return batch } // Use a map to keep only the last occurrence of each PK uniqueMap := make(map[interface{}][]interface{}) for _, row := range batch { pkValue := row[pkIndex] uniqueMap[pkValue] = row } // Convert back to slice result := make([][]interface{}, 0, len(uniqueMap)) for _, row := range uniqueMap { result = append(result, row) } return result } func processBatch( ctx context.Context, tx *sql.Tx, srcTableName string, dstTableName string, btch [][]interface{}, ) error { if len(btch) == 0 { return nil } batch := deduplicateBatch(btch, srcTableName) dstFields := bd.CdcTableFieldsMapping[srcTableName].DstFields fieldTypes := bd.CdcSrcTableFieldsTypes[dstTableName] pk := bd.CdcTableFieldsMapping[srcTableName].PrimaryKey // Build the base query baseQuery := fmt.Sprintf( `INSERT INTO %s (%s, synctime) VALUES `, dstTableName, strings.Join(dstFields, ", "), ) // Build the ON CONFLICT part var updateStatements []string for _, dstField := range dstFields { if dstField == pk { continue // Skip PK in update } updateStatements = append(updateStatements, fmt.Sprintf("%s = EXCLUDED.%s", dstField, dstField)) } updateStatements = append(updateStatements, "synctime = (CURRENT_TIMESTAMP AT TIME ZONE 'UTC')") conflictQuery := fmt.Sprintf( ` ON CONFLICT (%s) DO UPDATE SET %s`, pk, strings.Join(updateStatements, ", "), ) // Build the full query with all placeholders var placeholders []string var values []interface{} placeholderIndex := 1 for _, row := range batch { var rowPlaceholders []string for colIdx, dstField := range dstFields { fieldType, exists := fieldTypes[dstField] if !exists { return fmt.Errorf("field %s not found in CdcSrcTableFieldsTypes", dstField) } // Handle special types if needed switch { case strings.HasPrefix(fieldType, "varchar"), fieldType == "text": // For text-based fields, ensure proper escaping is handled by the driver rowPlaceholders = append(rowPlaceholders, fmt.Sprintf("$%d", placeholderIndex)) case strings.HasPrefix(fieldType, "int"): rowPlaceholders = append(rowPlaceholders, fmt.Sprintf("$%d", placeholderIndex)) case strings.HasPrefix(fieldType, "double"): rowPlaceholders = append(rowPlaceholders, fmt.Sprintf("$%d", placeholderIndex)) case strings.HasPrefix(fieldType, "timestamp"): rowPlaceholders = append(rowPlaceholders, fmt.Sprintf("$%d", placeholderIndex)) case strings.HasPrefix(fieldType, "numeric"): rowPlaceholders = append(rowPlaceholders, fmt.Sprintf("$%d", placeholderIndex)) default: return fmt.Errorf("unsupported field type: %s", fieldType) } values = append(values, row[colIdx]) placeholderIndex++ } rowPlaceholders = append(rowPlaceholders, "(CURRENT_TIMESTAMP AT TIME ZONE 'UTC')") placeholders = append(placeholders, "("+strings.Join(rowPlaceholders, ",")+")") } fullQuery := baseQuery + strings.Join(placeholders, ",") + conflictQuery // Execute the batch insert _, err := tx.ExecContext(ctx, fullQuery, values...) if err != nil { return fmt.Errorf("batch insert failed: %w. Query: %s", err, fullQuery) } return nil }