Skip to content

Commit 76aab08

Browse files
authored
Fix processing the same item twice for interrupted jobs (#13)
Fixes #11. Fixes #12.
1 parent 9bc75b7 commit 76aab08

9 files changed

Lines changed: 50 additions & 42 deletions

File tree

CHANGELOG.md

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,7 @@
11
## master (unreleased)
22

3+
- Fix processing the same item twice for interrupted jobs
4+
35
## 0.4.0 (2024-05-10)
46

57
- Support ordering using multiple directions for ActiveRecord enumerators

lib/sidekiq_iteration/active_record_enumerator.rb

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -78,7 +78,7 @@ def batches
7878
Enumerator.new(-> { records_size }) do |yielder|
7979
while (batch = next_batch(load: true))
8080
increment_iteration
81-
yielder.yield(batch, cursor_value(batch.last))
81+
yielder.yield(batch, cursor_value(batch.first))
8282
end
8383
end
8484
end

lib/sidekiq_iteration/iteration.rb

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -183,18 +183,19 @@ def iterate_with_enumerator(enumerator, arguments)
183183

184184
enumerator.each do |object_from_enumerator, index|
185185
found_record = true
186-
around_iteration do
187-
each_iteration(object_from_enumerator, *arguments)
188-
end
189186
@cursor_position = index
190-
@current_run_iterations += 1
191187

192188
throttle_condition = find_throttle_condition
193189
if throttle_condition
194190
@job_iteration_retry_backoff = throttle_condition.backoff
195191
@needs_reenqueue = true
196192
return false
197193
end
194+
195+
around_iteration do
196+
each_iteration(object_from_enumerator, *arguments)
197+
end
198+
@current_run_iterations += 1
198199
end
199200

200201
unless found_record

lib/sidekiq_iteration/nested_enumerator.rb

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,9 @@ def iterate(current_items, current_cursor, index, &block)
3333
else
3434
iterate(current_items + [item], current_cursor + [cursor_value], index + 1, &block)
3535
end
36+
37+
# Nullify downstream enums cursors on each iteration.
38+
@cursor.fill(nil, index + 1)
3639
end
3740
end
3841
end

test/active_record_enumerator_test.rb

Lines changed: 11 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -40,14 +40,14 @@ class ActiveRecordEnumeratorTest < TestCase
4040
end
4141
end
4242

43-
test "#batches yields batches of records with the last record's cursor position" do
43+
test "#batches yields batches of records with the first record's cursor position" do
4444
enum = build_enumerator.batches
4545
assert_equal(10, enum.size)
4646

4747
products = Product.order(:id).take(4).each_slice(2).to_a
4848

4949
enum.first(2).each_with_index do |(element, cursor), index|
50-
assert_equal([products[index], products[index].last.id], [element, cursor])
50+
assert_equal([products[index], products[index].first.id], [element, cursor])
5151
end
5252
end
5353

@@ -101,7 +101,7 @@ class ActiveRecordEnumeratorTest < TestCase
101101
enum = build_enumerator(batch_size: 4).batches
102102
products = Product.order(:id).take(4)
103103

104-
assert_equal([products, products.last.id], enum.first)
104+
assert_equal([products, products.first.id], enum.first)
105105
end
106106

107107
test "paginates using integer conditionals for primary key when no columns are defined" do
@@ -116,14 +116,14 @@ class ActiveRecordEnumeratorTest < TestCase
116116
enum = build_enumerator(columns: [:id]).batches
117117
products = Product.order(:id).take(2)
118118

119-
assert_equal([products, products.last.id], enum.first)
119+
assert_equal([products, products.first.id], enum.first)
120120
end
121121

122122
test "single column as symbol" do
123123
enum = build_enumerator(columns: :id).batches
124124
products = Product.order(:id).take(2)
125125

126-
assert_equal([products, products.last.id], enum.first)
126+
assert_equal([products, products.first.id], enum.first)
127127
end
128128

129129
test "raises when columns does not include primary key" do
@@ -142,7 +142,7 @@ class ActiveRecordEnumeratorTest < TestCase
142142
enum = build_enumerator(columns: [:updated_at, :id]).batches
143143
products = Product.order(:updated_at, :id).take(2)
144144

145-
assert_equal([products, [products.last.updated_at.strftime(SQL_TIME_FORMAT), products.last.id]], enum.first)
145+
assert_equal([products, [products.first.updated_at.strftime(SQL_TIME_FORMAT), products.first.id]], enum.first)
146146
end
147147

148148
test "columns configured with primary key only queries primary key column once" do
@@ -163,7 +163,7 @@ class ActiveRecordEnumeratorTest < TestCase
163163

164164
test ":order with single direction" do
165165
enum = build_enumerator(order: :desc).batches
166-
product_batches = Product.order(id: :desc).take(4).in_groups_of(2).map { |products| [products, products.last.id] }
166+
product_batches = Product.order(id: :desc).take(4).in_groups_of(2).map { |products| [products, products.first.id] }
167167

168168
enum.first(2).each_with_index do |(batch, cursor), index|
169169
assert_equal(product_batches[index], [batch, cursor])
@@ -203,24 +203,22 @@ class ActiveRecordEnumeratorTest < TestCase
203203
enum = build_enumerator(cursor: one.id).batches
204204

205205
queries = track_queries do
206-
assert_equal([[one, two], two.id], enum.first)
206+
assert_equal([[one, two], one.id], enum.first)
207207
end
208208
assert queries.any?(/"products"\."id" >= 1/)
209209

210-
assert_equal([[three, four], four.id], enum.first)
210+
assert_equal([[three, four], three.id], enum.first)
211211
end
212212

213213
test "can be resumed on multiple columns" do
214214
enum = build_enumerator(columns: [:created_at, :id]).batches
215215
products = Product.order(:created_at, :id).take(2)
216216

217-
cursor = [products.last.created_at.strftime(SQL_TIME_FORMAT), products.last.id]
217+
cursor = [products.first.created_at.strftime(SQL_TIME_FORMAT), products.first.id]
218218
assert_equal([products, cursor], enum.first)
219219

220220
enum = build_enumerator(columns: [:created_at, :id], cursor: cursor).batches
221-
products = Product.order(:created_at, :id).offset(1).take(2)
222-
223-
cursor = [products.last.created_at.strftime(SQL_TIME_FORMAT), products.last.id]
221+
cursor = [products.first.created_at.strftime(SQL_TIME_FORMAT), products.first.id]
224222
assert_equal([products, cursor], enum.first)
225223
end
226224

test/integration/integration_test.rb

Lines changed: 11 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -19,14 +19,14 @@ def setup
1919
start_worker_and_wait
2020

2121
assert_equal(1, queue_size)
22-
assert_equal(0, iteration_metadata["cursor_position"])
22+
assert_equal(1, iteration_metadata["cursor_position"])
2323
assert_equal(1, iteration_metadata["times_interrupted"])
2424

2525
conn.set("stop_after_num", 2)
2626
start_worker_and_wait
2727

2828
assert_equal(1, queue_size)
29-
assert_equal(2, iteration_metadata["cursor_position"])
29+
assert_equal(3, iteration_metadata["cursor_position"])
3030
assert_equal(2, iteration_metadata["times_interrupted"])
3131

3232
TerminateJob.perform_async
@@ -43,26 +43,27 @@ def setup
4343
start_worker_and_wait
4444

4545
Sidekiq.redis do |conn|
46-
assert_equal(2, iteration_metadata["cursor_position"])
46+
assert_equal(3, iteration_metadata["cursor_position"])
4747
assert_equal(1, iteration_metadata["executions"])
4848
assert_equal(0, iteration_metadata["times_interrupted"])
49-
assert_equal(3, conn.llen("records_performed"))
49+
assert_equal(["0", "1", "2"], conn.lrange("records_performed", 0, -1))
5050
assert_equal(1, conn.get("on_start_called").to_i)
5151

5252
start_worker_and_wait
5353

54-
assert_equal(4, iteration_metadata["cursor_position"])
54+
assert_equal(6, iteration_metadata["cursor_position"])
5555
assert_equal(2, iteration_metadata["executions"])
5656
assert_equal(0, iteration_metadata["times_interrupted"])
57-
assert_equal(6, conn.llen("records_performed"))
57+
assert_equal(["0", "1", "2", "3", "4", "5"], conn.lrange("records_performed", 0, -1))
5858
assert_equal(1, conn.get("on_start_called").to_i)
5959
assert_equal(0, conn.get("on_complete_called").to_i)
60-
end
6160

62-
# last attempt
63-
start_worker_and_wait
61+
# last attempt
62+
start_worker_and_wait
6463

65-
assert_equal(0, queue_size)
64+
assert_equal(1, Sidekiq::RetrySet.new.size)
65+
assert_equal(["0", "1", "2", "3", "4", "5", "6", "7", "8"], conn.lrange("records_performed", 0, -1))
66+
end
6667
end
6768

6869
test "failing non iteration job" do

test/integration/jobs.rb

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -25,7 +25,7 @@ class FailingIterationJob
2525
include SidekiqIteration::Iteration
2626

2727
sidekiq_options retry: 3
28-
sidekiq_retry_in { 1 }
28+
sidekiq_retry_in { 0 }
2929

3030
def on_start
3131
Sidekiq.redis do |conn|

test/iteration_test.rb

Lines changed: 15 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -140,7 +140,7 @@ def each_iteration(record)
140140
metadata = iteration_metadata(job)
141141
assert_equal(1, metadata["times_interrupted"])
142142
assert_equal(1, metadata["executions"])
143-
assert_equal(Product.second.id, metadata["cursor_position"])
143+
assert_equal(Product.third.id, metadata["cursor_position"])
144144

145145
ActiveRecordIterationJob.perform_one
146146
assert_equal(4, ActiveRecordIterationJob.records_performed.size)
@@ -188,7 +188,7 @@ def each_iteration(records)
188188

189189
job = peek_into_queue
190190
metadata = iteration_metadata(job)
191-
assert_equal(processed_records[2], metadata["cursor_position"])
191+
assert_equal(processed_records[3], metadata["cursor_position"])
192192
assert_equal(1, metadata["times_interrupted"])
193193
assert_equal(1, metadata["executions"])
194194

@@ -203,14 +203,13 @@ def each_iteration(records)
203203
metadata = iteration_metadata(job)
204204
assert_equal(2, metadata["times_interrupted"])
205205
assert_equal(2, metadata["executions"])
206-
assert_equal(processed_records[4], metadata["cursor_position"])
206+
assert_equal(processed_records[6], metadata["cursor_position"])
207207
continue_iterating(BatchActiveRecordIterationJob)
208208

209209
BatchActiveRecordIterationJob.perform_one
210210
assert_jobs_in_queue(0)
211211
assert_equal(4, BatchActiveRecordIterationJob.records_performed.size)
212-
# 10 records + 2 times restarted and ran on the same records
213-
assert_equal(12, BatchActiveRecordIterationJob.records_performed.flatten.size)
212+
assert_equal(processed_records.size, BatchActiveRecordIterationJob.records_performed.flatten.size)
214213

215214
assert_equal(1, BatchActiveRecordIterationJob.on_start_called)
216215
assert_equal(4, BatchActiveRecordIterationJob.around_iteration_called)
@@ -257,23 +256,27 @@ def each_iteration(record)
257256

258257
MultipleColumnsActiveRecordIterationJob.perform_async
259258

259+
products = Product.all.order(:updated_at, :id).to_a
260+
260261
1.upto(3) do |iter|
261262
MultipleColumnsActiveRecordIterationJob.perform_one
262263

263264
job = peek_into_queue
264265
last_processed_record = MultipleColumnsActiveRecordIterationJob.records_performed.last
265-
expected = [last_processed_record.updated_at.strftime("%Y-%m-%d %H:%M:%S.%6N"), last_processed_record.id]
266+
last_processed_record_index = products.index(last_processed_record)
267+
next_record = products[last_processed_record_index + 1]
268+
269+
expected = [next_record.updated_at.strftime("%Y-%m-%d %H:%M:%S.%6N"), next_record.id]
266270
metadata = iteration_metadata(job)
267271
assert_equal(expected, metadata["cursor_position"])
268272

269273
assert_equal(iter * 3, MultipleColumnsActiveRecordIterationJob.records_performed.size)
270274
end
271275

272-
products = Product.all.order(:updated_at, :id).to_a
273276
expected = [
274277
products[0], products[1], products[2], # first iteration
275-
products[2], products[3], products[4], # second iteration
276-
products[4], products[5], products[6] # second iteration
278+
products[3], products[4], products[5], # second iteration
279+
products[6], products[7], products[8] # second iteration
277280
]
278281
assert_equal(expected, MultipleColumnsActiveRecordIterationJob.records_performed)
279282
end
@@ -483,13 +486,13 @@ def each_iteration(*)
483486
iterate_exact_times(ArrayJob, 1) # starts from 1-st record
484487
ArrayJob.perform_inline
485488

486-
iterate_exact_times(ArrayJob, 3) # starts from 1-st record
489+
iterate_exact_times(ArrayJob, 3) # starts from 2-nd record
487490
ArrayJob.perform_one
488491

489-
continue_iterating(ArrayJob) # starts from 3-rd record
492+
continue_iterating(ArrayJob) # starts from 5-th record
490493
ArrayJob.perform_one
491494

492-
assert_equal([1, 3, 8], ArrayJob.current_run_iterations)
495+
assert_equal([1, 3, 6], ArrayJob.current_run_iterations)
493496
end
494497

495498
class IterateUntilMaxRuntimeJob < SimpleIterationJob

test/throttling_test.rb

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -53,7 +53,7 @@ def setup
5353
assert_equal(1, runs_metadata.fetch("cursor_position"))
5454
assert_equal(1, runs_metadata.fetch("times_interrupted"))
5555

56-
assert_equal([1, 2], ThrottleJob.iterations_performed)
56+
assert_equal([1], ThrottleJob.iterations_performed)
5757
end
5858

5959
test "pushes job back to the queue with the configured throttle backoff" do

0 commit comments

Comments
 (0)