This commit is contained in:
chrislu
2025-09-01 16:12:10 -07:00
parent 2fbceca959
commit e3798c2ec9
3 changed files with 79 additions and 44 deletions
+12 -8
View File
@@ -397,9 +397,6 @@ func convertHybridResultsToSQL(results []HybridScanResult, columns []string) *Qu
for columnName := range columnSet {
columns = append(columns, columnName)
}
// Add metadata columns showing data source
columns = append(columns, "_source")
}
// Convert to SQL rows
@@ -407,12 +404,19 @@ func convertHybridResultsToSQL(results []HybridScanResult, columns []string) *Qu
for i, result := range results {
row := make([]sqltypes.Value, len(columns))
for j, columnName := range columns {
if columnName == "_source" {
switch columnName {
case "_source":
row[j] = sqltypes.NewVarChar(result.Source)
} else if value, exists := result.Values[columnName]; exists {
row[j] = convertSchemaValueToSQL(value)
} else {
row[j] = sqltypes.NULL
case "_timestamp_ns":
row[j] = sqltypes.NewInt64(result.Timestamp)
case "_key":
row[j] = sqltypes.NewVarBinary(string(result.Key))
default:
if value, exists := result.Values[columnName]; exists {
row[j] = convertSchemaValueToSQL(value)
} else {
row[j] = sqltypes.NULL
}
}
}
rows[i] = row
+2 -3
View File
@@ -580,9 +580,6 @@ func (hms *HybridMessageScanner) ConvertToSQLResult(results []HybridScanResult,
for columnName := range columnSet {
columns = append(columns, columnName)
}
// Add metadata columns for debugging
columns = append(columns, "_source", "_timestamp_ns")
}
// Convert to SQL rows
@@ -595,6 +592,8 @@ func (hms *HybridMessageScanner) ConvertToSQLResult(results []HybridScanResult,
row[j] = sqltypes.NewVarChar(result.Source)
case "_timestamp_ns":
row[j] = sqltypes.NewInt64(result.Timestamp)
case "_key":
row[j] = sqltypes.NewVarBinary(string(result.Key))
default:
if value, exists := result.Values[columnName]; exists {
row[j] = convertSchemaValueToSQL(value)