Chunk SQL queries to avoid SQLite parameter limitsUnverified
381f2e4 parent: bb615b0 modified
internal/cluster/social.go +42 -16 | @@ -3,10 +3,25 @@ package cluster | ||
| 3 | 3 | import ( |
| 4 | 4 | "context" |
| 5 | 5 | "fmt" |
| 6 | + "iter" | |
| 6 | 7 | ) |
| 7 | 8 | |
| 8 | 9 | const maxFollowDepth = 3 |
| 9 | 10 | |
| 11 | +func chunk[T any](s []T, size int) iter.Seq[[]T] { | |
| 12 | + return func(yield func([]T) bool) { | |
| 13 | + for i := 0; i < len(s); i += size { | |
| 14 | + end := i + size | |
| 15 | + if end > len(s) { | |
| 16 | + end = len(s) | |
| 17 | + } | |
| 18 | + if !yield(s[i:end]) { | |
| 19 | + return | |
| 20 | + } | |
| 21 | + } | |
| 22 | + } | |
| 23 | +} | |
| 24 | + | |
| 10 | 25 | type followDistance struct { |
| 11 | 26 | userA string |
| 12 | 27 | userB string |
| @@ -125,17 +140,20 @@ func (e *Engine) ComputeFollowDistances(ctx context.Context) error { | ||
| 125 | 140 | } |
| 126 | 141 | defer func() { _ = tx.Rollback() }() |
| 127 | 142 | |
| 128 | - ph := make([]string, len(dirtyUsers)) | |
| 129 | - args := make([]any, len(dirtyUsers)) | |
| 130 | - for i, did := range dirtyUsers { | |
| 131 | - ph[i] = "?" | |
| 132 | - args[i] = did | |
| 133 | - } | |
| 134 | - if _, err := tx.ExecContext(ctx, | |
| 135 | - fmt.Sprintf("DELETE FROM recs.follow_distances WHERE user_a IN (%s)", joinPh(ph)), | |
| 136 | - args..., | |
| 137 | - ); err != nil { | |
| 138 | - return err | |
| 143 | + const sqliteMaxVars = 500 | |
| 144 | + for chunk := range chunk(dirtyUsers, sqliteMaxVars) { | |
| 145 | + ph := make([]string, len(chunk)) | |
| 146 | + args := make([]any, len(chunk)) | |
| 147 | + for i, did := range chunk { | |
| 148 | + ph[i] = "?" | |
| 149 | + args[i] = did | |
| 150 | + } | |
| 151 | + if _, err := tx.ExecContext(ctx, | |
| 152 | + fmt.Sprintf("DELETE FROM recs.follow_distances WHERE user_a IN (%s)", joinPh(ph)), | |
| 153 | + args..., | |
| 154 | + ); err != nil { | |
| 155 | + return err | |
| 156 | + } | |
| 139 | 157 | } |
| 140 | 158 | |
| 141 | 159 | stmt, err := tx.PrepareContext(ctx, `INSERT INTO recs.follow_distances (user_a, user_b, distance) VALUES (?, ?, ?)`) |
| @@ -150,11 +168,19 @@ func (e *Engine) ComputeFollowDistances(ctx context.Context) error { | ||
| 150 | 168 | } |
| 151 | 169 | } |
| 152 | 170 | |
| 153 | - if _, err := tx.ExecContext(ctx, | |
| 154 | - fmt.Sprintf("UPDATE main.users SET follows_dirty = 0 WHERE did IN (%s)", joinPh(ph)), | |
| 155 | - args..., | |
| 156 | - ); err != nil { | |
| 157 | - return err | |
| 171 | + for chunk := range chunk(dirtyUsers, sqliteMaxVars) { | |
| 172 | + ph := make([]string, len(chunk)) | |
| 173 | + args := make([]any, len(chunk)) | |
| 174 | + for i, did := range chunk { | |
| 175 | + ph[i] = "?" | |
| 176 | + args[i] = did | |
| 177 | + } | |
| 178 | + if _, err := tx.ExecContext(ctx, | |
| 179 | + fmt.Sprintf("UPDATE main.users SET follows_dirty = 0 WHERE did IN (%s)", joinPh(ph)), | |
| 180 | + args..., | |
| 181 | + ); err != nil { | |
| 182 | + return err | |
| 183 | + } | |
| 158 | 184 | } |
| 159 | 185 | |
| 160 | 186 | e.logger.Info("follow distances computed", "users", len(dirtyUsers), "pairs", len(distances)) |
| @@ -3,10 +3,25 @@ package cluster | |||
| 3 | import ( | 3 | import ( |
| 4 | "context" | 4 | "context" |
| 5 | "fmt" | 5 | "fmt" |
| 6 | + "iter" | ||
| 6 | ) | 7 | ) |
| 7 | 8 | ||
| 8 | const maxFollowDepth = 3 | 9 | const maxFollowDepth = 3 |
| 9 | 10 | ||
| 11 | +func chunk[T any](s []T, size int) iter.Seq[[]T] { | ||
| 12 | + return func(yield func([]T) bool) { | ||
| 13 | + for i := 0; i < len(s); i += size { | ||
| 14 | + end := i + size | ||
| 15 | + if end > len(s) { | ||
| 16 | + end = len(s) | ||
| 17 | + } | ||
| 18 | + if !yield(s[i:end]) { | ||
| 19 | + return | ||
| 20 | + } | ||
| 21 | + } | ||
| 22 | + } | ||
| 23 | +} | ||
| 24 | + | ||
| 10 | type followDistance struct { | 25 | type followDistance struct { |
| 11 | userA string | 26 | userA string |
| 12 | userB string | 27 | userB string |
| @@ -125,17 +140,20 @@ func (e *Engine) ComputeFollowDistances(ctx context.Context) error { | |||
| 125 | } | 140 | } |
| 126 | defer func() { _ = tx.Rollback() }() | 141 | defer func() { _ = tx.Rollback() }() |
| 127 | 142 | ||
| 128 | - ph := make([]string, len(dirtyUsers)) | 143 | + const sqliteMaxVars = 500 |
| 129 | - args := make([]any, len(dirtyUsers)) | 144 | + for chunk := range chunk(dirtyUsers, sqliteMaxVars) { |
| 130 | - for i, did := range dirtyUsers { | 145 | + ph := make([]string, len(chunk)) |
| 131 | - ph[i] = "?" | 146 | + args := make([]any, len(chunk)) |
| 132 | - args[i] = did | 147 | + for i, did := range chunk { |
| 133 | - } | 148 | + ph[i] = "?" |
| 134 | - if _, err := tx.ExecContext(ctx, | 149 | + args[i] = did |
| 135 | - fmt.Sprintf("DELETE FROM recs.follow_distances WHERE user_a IN (%s)", joinPh(ph)), | 150 | + } |
| 136 | - args..., | 151 | + if _, err := tx.ExecContext(ctx, |
| 137 | - ); err != nil { | 152 | + fmt.Sprintf("DELETE FROM recs.follow_distances WHERE user_a IN (%s)", joinPh(ph)), |
| 138 | - return err | 153 | + args..., |
| 154 | + ); err != nil { | ||
| 155 | + return err | ||
| 156 | + } | ||
| 139 | } | 157 | } |
| 140 | 158 | ||
| 141 | stmt, err := tx.PrepareContext(ctx, `INSERT INTO recs.follow_distances (user_a, user_b, distance) VALUES (?, ?, ?)`) | 159 | stmt, err := tx.PrepareContext(ctx, `INSERT INTO recs.follow_distances (user_a, user_b, distance) VALUES (?, ?, ?)`) |
| @@ -150,11 +168,19 @@ func (e *Engine) ComputeFollowDistances(ctx context.Context) error { | |||
| 150 | } | 168 | } |
| 151 | } | 169 | } |
| 152 | 170 | ||
| 153 | - if _, err := tx.ExecContext(ctx, | 171 | + for chunk := range chunk(dirtyUsers, sqliteMaxVars) { |
| 154 | - fmt.Sprintf("UPDATE main.users SET follows_dirty = 0 WHERE did IN (%s)", joinPh(ph)), | 172 | + ph := make([]string, len(chunk)) |
| 155 | - args..., | 173 | + args := make([]any, len(chunk)) |
| 156 | - ); err != nil { | 174 | + for i, did := range chunk { |
| 157 | - return err | 175 | + ph[i] = "?" |
| 176 | + args[i] = did | ||
| 177 | + } | ||
| 178 | + if _, err := tx.ExecContext(ctx, | ||
| 179 | + fmt.Sprintf("UPDATE main.users SET follows_dirty = 0 WHERE did IN (%s)", joinPh(ph)), | ||
| 180 | + args..., | ||
| 181 | + ); err != nil { | ||
| 182 | + return err | ||
| 183 | + } | ||
| 158 | } | 184 | } |
| 159 | 185 | ||
| 160 | e.logger.Info("follow distances computed", "users", len(dirtyUsers), "pairs", len(distances)) | 186 | e.logger.Info("follow distances computed", "users", len(dirtyUsers), "pairs", len(distances)) |