Skip to content

Commit 53b9faf

Browse files
committed
Fix builder's chunk order and drop panic-on-warnings row count check
- Removed count subqueries from range builders to restore the more performant approach from PR github#471 - Modified panic-on-warnings logic to trigger errors based solely on SQL warnings, not on row count mismatches - This addresses potential race conditions where row count comparisons could produce false positives due to concurrent table modifications
1 parent c4ac9ba commit 53b9faf

5 files changed

Lines changed: 128 additions & 77 deletions

File tree

go/logic/applier.go

Lines changed: 38 additions & 36 deletions
Original file line numberDiff line numberDiff line change
@@ -663,47 +663,49 @@ func (this *Applier) ReadMigrationRangeValues() error {
663663
// which will be used for copying the next chunk of rows. Ir returns "false" if there is
664664
// no further chunk to work through, i.e. we're past the last chunk and are done with
665665
// iterating the range (and this done with copying row chunks)
666-
func (this *Applier) CalculateNextIterationRangeEndValues() (hasFurtherRange bool, expectedRowCount int64, err error) {
667-
query, explodedArgs, err := sql.BuildUniqueKeyRangeEndPreparedQueryViaTemptable(
668-
this.migrationContext.DatabaseName,
669-
this.migrationContext.OriginalTableName,
670-
&this.migrationContext.UniqueKey.Columns,
671-
this.migrationContext.MigrationIterationRangeMinValues.AbstractValues(),
672-
this.migrationContext.MigrationRangeMaxValues.AbstractValues(),
673-
atomic.LoadInt64(&this.migrationContext.ChunkSize),
674-
this.migrationContext.GetIteration() == 0,
675-
fmt.Sprintf("iteration:%d", this.migrationContext.GetIteration()),
676-
)
677-
if err != nil {
678-
return hasFurtherRange, expectedRowCount, err
679-
}
680-
681-
rows, err := this.db.Query(query, explodedArgs...)
682-
if err != nil {
683-
return hasFurtherRange, expectedRowCount, err
684-
}
685-
defer rows.Close()
686-
687-
iterationRangeMaxValues := sql.NewColumnValues(this.migrationContext.UniqueKey.Len() + 1)
688-
for rows.Next() {
689-
if err = rows.Scan(iterationRangeMaxValues.ValuesPointers...); err != nil {
690-
return hasFurtherRange, expectedRowCount, err
666+
func (this *Applier) CalculateNextIterationRangeEndValues() (hasFurtherRange bool, err error) {
667+
for i := 0; i < 2; i++ {
668+
buildFunc := sql.BuildUniqueKeyRangeEndPreparedQueryViaOffset
669+
if i == 1 {
670+
buildFunc = sql.BuildUniqueKeyRangeEndPreparedQueryViaTemptable
671+
}
672+
query, explodedArgs, err := buildFunc(
673+
this.migrationContext.DatabaseName,
674+
this.migrationContext.OriginalTableName,
675+
&this.migrationContext.UniqueKey.Columns,
676+
this.migrationContext.MigrationIterationRangeMinValues.AbstractValues(),
677+
this.migrationContext.MigrationRangeMaxValues.AbstractValues(),
678+
atomic.LoadInt64(&this.migrationContext.ChunkSize),
679+
this.migrationContext.GetIteration() == 0,
680+
fmt.Sprintf("iteration:%d", this.migrationContext.GetIteration()),
681+
)
682+
if err != nil {
683+
return hasFurtherRange, err
691684
}
692685

693-
expectedRowCount = (*iterationRangeMaxValues.ValuesPointers[len(iterationRangeMaxValues.ValuesPointers)-1].(*interface{})).(int64)
694-
iterationRangeMaxValues = sql.ToColumnValues(iterationRangeMaxValues.AbstractValues()[:len(iterationRangeMaxValues.AbstractValues())-1])
686+
rows, err := this.db.Query(query, explodedArgs...)
687+
if err != nil {
688+
return hasFurtherRange, err
689+
}
690+
defer rows.Close()
695691

696-
hasFurtherRange = expectedRowCount > 0
697-
}
698-
if err = rows.Err(); err != nil {
699-
return hasFurtherRange, expectedRowCount, err
700-
}
701-
if hasFurtherRange {
702-
this.migrationContext.MigrationIterationRangeMaxValues = iterationRangeMaxValues
703-
return hasFurtherRange, expectedRowCount, nil
692+
iterationRangeMaxValues := sql.NewColumnValues(this.migrationContext.UniqueKey.Len())
693+
for rows.Next() {
694+
if err = rows.Scan(iterationRangeMaxValues.ValuesPointers...); err != nil {
695+
return hasFurtherRange, err
696+
}
697+
hasFurtherRange = true
698+
}
699+
if err = rows.Err(); err != nil {
700+
return hasFurtherRange, err
701+
}
702+
if hasFurtherRange {
703+
this.migrationContext.MigrationIterationRangeMaxValues = iterationRangeMaxValues
704+
return hasFurtherRange, nil
705+
}
704706
}
705707
this.migrationContext.Log.Debugf("Iteration complete: no further range to iterate")
706-
return hasFurtherRange, expectedRowCount, nil
708+
return hasFurtherRange, nil
707709
}
708710

709711
// ApplyIterationInsertQuery issues a chunk-INSERT query on the ghost table. It is where

go/logic/applier_test.go

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -563,10 +563,9 @@ func (suite *ApplierTestSuite) TestPanicOnWarningsInApplyIterationInsertQuerySuc
563563
suite.Require().NoError(err)
564564

565565
migrationContext.SetNextIterationRangeMinValues()
566-
hasFurtherRange, expectedRangeSize, err := applier.CalculateNextIterationRangeEndValues()
566+
hasFurtherRange, err := applier.CalculateNextIterationRangeEndValues()
567567
suite.Require().NoError(err)
568568
suite.Require().True(hasFurtherRange)
569-
suite.Require().Equal(int64(1), expectedRangeSize)
570569

571570
_, rowsAffected, _, err := applier.ApplyIterationInsertQuery()
572571
suite.Require().NoError(err)

go/logic/migrator.go

Lines changed: 3 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1344,7 +1344,7 @@ func (this *Migrator) iterateChunks() error {
13441344
}
13451345

13461346
// When hasFurtherRange is false, original table might be write locked and CalculateNextIterationRangeEndValues would hangs forever
1347-
hasFurtherRange, expectedRangeSize, err := this.applier.CalculateNextIterationRangeEndValues()
1347+
hasFurtherRange, err := this.applier.CalculateNextIterationRangeEndValues()
13481348
if err != nil {
13491349
return err // wrapping call will retry
13501350
}
@@ -1373,10 +1373,8 @@ func (this *Migrator) iterateChunks() error {
13731373
for _, warning := range this.migrationContext.MigrationLastInsertSQLWarnings {
13741374
this.migrationContext.Log.Infof("ApplyIterationInsertQuery has SQL warnings! %s", warning)
13751375
}
1376-
if expectedRangeSize != rowsAffected {
1377-
joinedWarnings := strings.Join(this.migrationContext.MigrationLastInsertSQLWarnings, "; ")
1378-
terminateRowIteration(fmt.Errorf("ApplyIterationInsertQuery failed because of SQL warnings: [%s]", joinedWarnings))
1379-
}
1376+
joinedWarnings := strings.Join(this.migrationContext.MigrationLastInsertSQLWarnings, "; ")
1377+
terminateRowIteration(fmt.Errorf("ApplyIterationInsertQuery failed because of SQL warnings: [%s]", joinedWarnings))
13801378
}
13811379
}
13821380

go/sql/builder.go

Lines changed: 57 additions & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -251,6 +251,59 @@ func BuildRangeInsertPreparedQuery(databaseName, originalTableName, ghostTableNa
251251
return BuildRangeInsertQuery(databaseName, originalTableName, ghostTableName, sharedColumns, mappedSharedColumns, uniqueKey, uniqueKeyColumns, rangeStartValues, rangeEndValues, rangeStartArgs, rangeEndArgs, includeRangeStartValues, transactionalTable, noWait)
252252
}
253253

254+
func BuildUniqueKeyRangeEndPreparedQueryViaOffset(databaseName, tableName string, uniqueKeyColumns *ColumnList, rangeStartArgs, rangeEndArgs []interface{}, chunkSize int64, includeRangeStartValues bool, hint string) (result string, explodedArgs []interface{}, err error) {
255+
if uniqueKeyColumns.Len() == 0 {
256+
return "", explodedArgs, fmt.Errorf("Got 0 columns in BuildUniqueKeyRangeEndPreparedQuery")
257+
}
258+
databaseName = EscapeName(databaseName)
259+
tableName = EscapeName(tableName)
260+
261+
var startRangeComparisonSign ValueComparisonSign = GreaterThanComparisonSign
262+
if includeRangeStartValues {
263+
startRangeComparisonSign = GreaterThanOrEqualsComparisonSign
264+
}
265+
rangeStartComparison, rangeExplodedArgs, err := BuildRangePreparedComparison(uniqueKeyColumns, rangeStartArgs, startRangeComparisonSign)
266+
if err != nil {
267+
return "", explodedArgs, err
268+
}
269+
explodedArgs = append(explodedArgs, rangeExplodedArgs...)
270+
rangeEndComparison, rangeExplodedArgs, err := BuildRangePreparedComparison(uniqueKeyColumns, rangeEndArgs, LessThanOrEqualsComparisonSign)
271+
if err != nil {
272+
return "", explodedArgs, err
273+
}
274+
explodedArgs = append(explodedArgs, rangeExplodedArgs...)
275+
276+
uniqueKeyColumnNames := duplicateNames(uniqueKeyColumns.Names())
277+
uniqueKeyColumnAscending := make([]string, len(uniqueKeyColumnNames))
278+
for i, column := range uniqueKeyColumns.Columns() {
279+
uniqueKeyColumnNames[i] = EscapeName(uniqueKeyColumnNames[i])
280+
if column.Type == EnumColumnType {
281+
uniqueKeyColumnAscending[i] = fmt.Sprintf("concat(%s) asc", uniqueKeyColumnNames[i])
282+
} else {
283+
uniqueKeyColumnAscending[i] = fmt.Sprintf("%s asc", uniqueKeyColumnNames[i])
284+
}
285+
}
286+
result = fmt.Sprintf(`
287+
select /* gh-ost %s.%s %s */
288+
%s
289+
from
290+
%s.%s
291+
where
292+
%s and %s
293+
order by
294+
%s
295+
limit 1
296+
offset %d`,
297+
databaseName, tableName, hint,
298+
strings.Join(uniqueKeyColumnNames, ", "),
299+
databaseName, tableName,
300+
rangeStartComparison, rangeEndComparison,
301+
strings.Join(uniqueKeyColumnAscending, ", "),
302+
(chunkSize - 1),
303+
)
304+
return result, explodedArgs, nil
305+
}
306+
254307
func BuildUniqueKeyRangeEndPreparedQueryViaTemptable(databaseName, tableName string, uniqueKeyColumns *ColumnList, rangeStartArgs, rangeEndArgs []interface{}, chunkSize int64, includeRangeStartValues bool, hint string) (result string, explodedArgs []interface{}, err error) {
255308
if uniqueKeyColumns.Len() == 0 {
256309
return "", explodedArgs, fmt.Errorf("Got 0 columns in BuildUniqueKeyRangeEndPreparedQuery")
@@ -286,22 +339,8 @@ func BuildUniqueKeyRangeEndPreparedQueryViaTemptable(databaseName, tableName str
286339
uniqueKeyColumnDescending[i] = fmt.Sprintf("%s desc", uniqueKeyColumnNames[i])
287340
}
288341
}
289-
290-
joinedColumnNames := strings.Join(uniqueKeyColumnNames, ", ")
291342
result = fmt.Sprintf(`
292-
select /* gh-ost %s.%s %s */
293-
%s,
294-
(select count(*) from (
295-
select
296-
%s
297-
from
298-
%s.%s
299-
where
300-
%s and %s
301-
order by
302-
%s
303-
limit %d
304-
) select_osc_chunk)
343+
select /* gh-ost %s.%s %s */ %s
305344
from (
306345
select
307346
%s
@@ -316,17 +355,13 @@ func BuildUniqueKeyRangeEndPreparedQueryViaTemptable(databaseName, tableName str
316355
order by
317356
%s
318357
limit 1`,
319-
databaseName, tableName, hint, joinedColumnNames,
320-
joinedColumnNames, databaseName, tableName,
321-
rangeStartComparison, rangeEndComparison,
322-
strings.Join(uniqueKeyColumnAscending, ", "), chunkSize,
323-
joinedColumnNames, databaseName, tableName,
358+
databaseName, tableName, hint, strings.Join(uniqueKeyColumnNames, ", "),
359+
strings.Join(uniqueKeyColumnNames, ", "), databaseName, tableName,
324360
rangeStartComparison, rangeEndComparison,
325361
strings.Join(uniqueKeyColumnAscending, ", "), chunkSize,
326362
strings.Join(uniqueKeyColumnDescending, ", "),
327363
)
328-
// 2x the explodedArgs for the subquery (CTE would be possible but not supported by MySQL 5)
329-
return result, append(explodedArgs, explodedArgs...), nil
364+
return result, explodedArgs, nil
330365
}
331366

332367
func BuildUniqueKeyMinValuesPreparedQuery(databaseName, tableName string, uniqueKey *UniqueKey) (string, error) {

go/sql/builder_test.go

Lines changed: 29 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -325,6 +325,33 @@ func TestBuildRangeInsertPreparedQuery(t *testing.T) {
325325
}
326326
}
327327

328+
func TestBuildUniqueKeyRangeEndPreparedQueryViaOffset(t *testing.T) {
329+
databaseName := "mydb"
330+
originalTableName := "tbl"
331+
var chunkSize int64 = 500
332+
{
333+
uniqueKeyColumns := NewColumnList([]string{"name", "position"})
334+
rangeStartArgs := []interface{}{3, 17}
335+
rangeEndArgs := []interface{}{103, 117}
336+
337+
query, explodedArgs, err := BuildUniqueKeyRangeEndPreparedQueryViaOffset(databaseName, originalTableName, uniqueKeyColumns, rangeStartArgs, rangeEndArgs, chunkSize, false, "test")
338+
require.NoError(t, err)
339+
expected := `
340+
select /* gh-ost mydb.tbl test */
341+
name, position
342+
from
343+
mydb.tbl
344+
where
345+
((name > ?) or (((name = ?)) AND (position > ?))) and ((name < ?) or (((name = ?)) AND (position < ?)) or ((name = ?) and (position = ?)))
346+
order by
347+
name asc, position asc
348+
limit 1
349+
offset 499`
350+
require.Equal(t, normalizeQuery(expected), normalizeQuery(query))
351+
require.Equal(t, []interface{}{3, 3, 17, 103, 103, 117, 103, 117}, explodedArgs)
352+
}
353+
}
354+
328355
func TestBuildUniqueKeyRangeEndPreparedQueryViaTemptable(t *testing.T) {
329356
databaseName := "mydb"
330357
originalTableName := "tbl"
@@ -338,17 +365,7 @@ func TestBuildUniqueKeyRangeEndPreparedQueryViaTemptable(t *testing.T) {
338365
require.NoError(t, err)
339366
expected := `
340367
select /* gh-ost mydb.tbl test */
341-
name, position,
342-
(select count(*) from (
343-
select
344-
name, position
345-
from
346-
mydb.tbl
347-
where ((name > ?) or (((name = ?)) AND (position > ?))) and ((name < ?) or (((name = ?)) AND (position < ?)) or ((name = ?) and (position = ?)))
348-
order by
349-
name asc, position asc
350-
limit 500
351-
) select_osc_chunk)
368+
name, position
352369
from (
353370
select
354371
name, position
@@ -363,7 +380,7 @@ func TestBuildUniqueKeyRangeEndPreparedQueryViaTemptable(t *testing.T) {
363380
name desc, position desc
364381
limit 1`
365382
require.Equal(t, normalizeQuery(expected), normalizeQuery(query))
366-
require.Equal(t, []interface{}{3, 3, 17, 103, 103, 117, 103, 117, 3, 3, 17, 103, 103, 117, 103, 117}, explodedArgs)
383+
require.Equal(t, []interface{}{3, 3, 17, 103, 103, 117, 103, 117}, explodedArgs)
367384
}
368385
}
369386

0 commit comments

Comments
 (0)