Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 7 additions & 0 deletions internal/postgres/retrier/pg_querier_retrier.go
Original file line number Diff line number Diff line change
Expand Up @@ -152,6 +152,13 @@ func (q *Querier) resetConn(ctx context.Context) error {
}

func (q *Querier) isRetriableError(err error) bool {
return IsRetriableError(err)
}

// IsRetriableError reports whether retrying an operation that failed with err
// could succeed. Callers retrying at a coarser granularity than a single
// query use it so their rule cannot drift from this one.
func IsRetriableError(err error) bool {
mappedErr := postgres.MapError(err)

permissionDenied := &postgres.ErrPermissionDenied{}
Expand Down
7 changes: 7 additions & 0 deletions internal/sync/semaphore.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,3 +17,10 @@ type WeightedSemaphore interface {
func NewWeightedSemaphore(size int64) *semaphore.Weighted {
return semaphore.NewWeighted(size)
}

// CopyBudgetReserve leaves room for non-copy connections
const CopyBudgetReserve = 5

func CopyBudgetSize(maxConnections int32) int64 {
return max(1, int64(maxConnections)-CopyBudgetReserve)
}
18 changes: 18 additions & 0 deletions internal/sync/semaphore_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,18 @@
// SPDX-License-Identifier: Apache-2.0

package sync

import (
"testing"

"github.com/stretchr/testify/require"
)

func TestCopyBudgetSize(t *testing.T) {
t.Parallel()

// zero would block every copy forever
require.Equal(t, int64(1), CopyBudgetSize(1))
require.Equal(t, int64(1), CopyBudgetSize(CopyBudgetReserve))
require.Equal(t, int64(45), CopyBudgetSize(50))
}
3 changes: 3 additions & 0 deletions pkg/wal/processor/postgres/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,9 @@ func (c *Config) retryPolicy() backoff.Config {
}
}

// EffectiveRetryPolicy returns the retry policy once the default is applied.
func (c *Config) EffectiveRetryPolicy() backoff.Config { return c.retryPolicy() }

func (c *Config) poolOptions() []pglib.PoolOption {
if c.MaxConnections == 0 {
return nil
Expand Down
11 changes: 2 additions & 9 deletions pkg/wal/processor/postgres/postgres_bulk_ingest_writer.go
Original file line number Diff line number Diff line change
Expand Up @@ -27,15 +27,12 @@ type BulkIngestWriter struct {
// copyBudget caps the total number of concurrent COPYs across all tables
// (and all their send drainers) so they never exhaust the target
// connection pool. It is sized from the resolved pool max-connections
// value, minus copyBudgetReserve.
// value, minus synclib.CopyBudgetReserve.
copyBudget synclib.WeightedSemaphore
}

const bulkIngestWriter = "postgres_bulk_ingest_writer"

// batch writer and retrier reset share this pool
const copyBudgetReserve = 5

var errUnexpectedCopiedRows = errors.New("number of rows copied doesn't match the source rows")

// NewBulkIngestWriter returns a postgres processor that batches and writes data
Expand All @@ -54,7 +51,7 @@ func NewBulkIngestWriter(ctx context.Context, config *Config, opts ...WriterOpti
biw := &BulkIngestWriter{
Writer: w,
batchSenderMap: synclib.NewMap[string, queryBatchSender](),
copyBudget: synclib.NewWeightedSemaphore(copyBudgetSize(w.maxConnections)),
copyBudget: synclib.NewWeightedSemaphore(synclib.CopyBudgetSize(w.maxConnections)),
}

biw.batchSenderBuilder = func(ctx context.Context, schema, table string) (queryBatchSender, error) {
Expand All @@ -68,10 +65,6 @@ func NewBulkIngestWriter(ctx context.Context, config *Config, opts ...WriterOpti
return biw, nil
}

func copyBudgetSize(maxConnections int32) int64 {
return max(1, int64(maxConnections)-copyBudgetReserve)
}

// ProcessWALEvent is called on every new message from the wal. It can be called
// concurrently.
func (w *BulkIngestWriter) ProcessWALEvent(ctx context.Context, walEvent *wal.Event) (err error) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -424,7 +424,7 @@ func TestBulkIngestWriter_sendBatch(t *testing.T) {
pgConn: tc.pgConn,
disableTriggers: tc.disableTriggers,
},
copyBudget: synclib.NewWeightedSemaphore(pglib.MaxConns - copyBudgetReserve),
copyBudget: synclib.NewWeightedSemaphore(pglib.MaxConns - synclib.CopyBudgetReserve),
}

err := writer.sendBatch(context.Background(), tc.batch)
Expand Down Expand Up @@ -508,14 +508,14 @@ func TestCopyBudgetSize(t *testing.T) {
}{
{name: "default pool", maxConnections: pglib.MaxConns, expected: 45},
{name: "configured pool", maxConnections: 12, expected: 7},
{name: "reserve matches pool", maxConnections: copyBudgetReserve, expected: 1},
{name: "reserve matches pool", maxConnections: synclib.CopyBudgetReserve, expected: 1},
{name: "pool smaller than reserve", maxConnections: 2, expected: 1},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
require.Equal(t, tt.expected, copyBudgetSize(tt.maxConnections))
require.Equal(t, tt.expected, synclib.CopyBudgetSize(tt.maxConnections))
})
}
}
Expand Down Expand Up @@ -579,7 +579,7 @@ func TestNewBulkIngestWriter_maxConnections(t *testing.T) {
require.True(t, ok)
require.Equal(t, tt.expectedObserver, observerPool.Config().MaxConns)

budget := copyBudgetSize(tt.expected)
budget := synclib.CopyBudgetSize(tt.expected)
for range budget {
require.True(t, writer.copyBudget.TryAcquire(1))
}
Expand Down