diff --git a/lib/benchmark/append_no_fail_if.rb b/lib/benchmark/append_no_fail_if.rb index bf149fe..e489e5d 100644 --- a/lib/benchmark/append_no_fail_if.rb +++ b/lib/benchmark/append_no_fail_if.rb @@ -15,9 +15,10 @@ module En57 end call do |measure, run_id| - type = "event_benchmarked" - tags = %W[writer:#{run_id}] - events = Array.new(@batch_size) { Event.new(type: type, tags: tags) } + events = + @batch_size.times.map do + Event.new(type: "event_benchmarked", tags: ["writer:#{run_id}"]) + end measure.call { @event_store.append(events) } end diff --git a/lib/benchmark/append_non_conflicting_tags.rb b/lib/benchmark/append_non_conflicting_tags.rb index 070bb50..40e23b7 100644 --- a/lib/benchmark/append_non_conflicting_tags.rb +++ b/lib/benchmark/append_non_conflicting_tags.rb @@ -15,10 +15,12 @@ module En57 end call do |measure, run_id| - type = "event_benchmarked" - tags = %W[writer:#{run_id}] - scope = @event_store.read.of_type(type).with_tag(tags) - events = Array.new(@batch_size) { Event.new(type: type, tags: tags) } + scope = + @event_store + .read + .of_type(type = "event_benchmarked") + .with_tag(tags = ["writer:#{run_id}"]) + events = @batch_size.times.map { Event.new(type:, tags:) } measure.call do begin diff --git a/lib/benchmark/concurrent_append_conflicting_tags.rb b/lib/benchmark/concurrent_append_conflicting_tags.rb index fae0783..63de037 100644 --- a/lib/benchmark/concurrent_append_conflicting_tags.rb +++ b/lib/benchmark/concurrent_append_conflicting_tags.rb @@ -14,18 +14,18 @@ module En57 end call do |measure, run_id| - type = "event_benchmarked" - tags = %W[writer:#{run_id}] concurrently do |_writer_id, barrier| - scope = @event_store.read.of_type(type).with_tag(tags) - events = Array.new(@batch_size) { Event.new(type: type, tags: tags) } - position = 0 - + scope = + @event_store + .read + .of_type(type = "event_benchmarked") + .with_tag(tags = ["writer:#{run_id}"]) + events = @batch_size.times.map { Event.new(type:, tags:) } barrier.wait measure.call do begin - @event_store.append(events, fail_if: scope.after(position)) + @event_store.append(events, fail_if: scope.after(position = 0)) rescue AppendConditionViolated record_retry scope.each_with_position do |_event, event_position| diff --git a/lib/benchmark/concurrent_append_no_fail_if.rb b/lib/benchmark/concurrent_append_no_fail_if.rb index 93fbd75..48ba19f 100644 --- a/lib/benchmark/concurrent_append_no_fail_if.rb +++ b/lib/benchmark/concurrent_append_no_fail_if.rb @@ -14,11 +14,14 @@ module En57 end call do |measure, _run_id| - type = "event_benchmarked" concurrently do |writer_id, barrier| - tags = %W[writer:#{writer_id}] - events = Array.new(@batch_size) { Event.new(type: type, tags: tags) } - + events = + @batch_size.times.map do + Event.new( + type: "event_benchmarked", + tags: ["writer:#{writer_id}"], + ) + end barrier.wait measure.call { @event_store.append(events) } diff --git a/lib/benchmark/concurrent_append_non_conflicting_tags.rb b/lib/benchmark/concurrent_append_non_conflicting_tags.rb index 7664e42..0e846c9 100644 --- a/lib/benchmark/concurrent_append_non_conflicting_tags.rb +++ b/lib/benchmark/concurrent_append_non_conflicting_tags.rb @@ -14,12 +14,13 @@ module En57 end call do |measure, _run_id| - type = "event_benchmarked" concurrently do |writer_id, barrier| - tags = %W[writer:#{writer_id}] - scope = @event_store.read.of_type(type).with_tag(tags) - events = Array.new(@batch_size) { Event.new(type: type, tags: tags) } - + scope = + @event_store + .read + .of_type(type = "event_benchmarked") + .with_tag(tags = ["writer:#{writer_id}"]) + events = @batch_size.times.map { Event.new(type:, tags:) } barrier.wait measure.call do diff --git a/lib/benchmark/concurrent_append_non_conflicting_tags_seeded.rb b/lib/benchmark/concurrent_append_non_conflicting_tags_seeded.rb index 3510835..ff4b599 100644 --- a/lib/benchmark/concurrent_append_non_conflicting_tags_seeded.rb +++ b/lib/benchmark/concurrent_append_non_conflicting_tags_seeded.rb @@ -14,12 +14,13 @@ module En57 end call do |measure, _run_id| - type = "event_benchmarked" concurrently do |writer_id, barrier| - tags = %W[writer:#{writer_id}] - scope = @event_store.read.of_type(type).with_tag(tags) - events = Array.new(@batch_size) { Event.new(type: type, tags: tags) } - + scope = + @event_store + .read + .of_type(type = "event_benchmarked") + .with_tag(tags = ["writer:#{writer_id}"]) + events = @batch_size.times.map { Event.new(type:, tags:) } barrier.wait measure.call do diff --git a/lib/benchmark/res_append_stream_any.rb b/lib/benchmark/res_append_stream_any.rb index 88b0217..b7ec6cb 100644 --- a/lib/benchmark/res_append_stream_any.rb +++ b/lib/benchmark/res_append_stream_any.rb @@ -1,8 +1,5 @@ # frozen_string_literal: true -require "active_record" -require "rails_event_store" - module En57 module Benchmark Scenario.define do @@ -13,20 +10,29 @@ module En57 batch_size 100 setup do |database_url| + require "active_record" + require "rails_event_store" + ActiveRecord::Base.establish_connection(database_url) @event_store = RailsEventStore::JSONClient.new end call do |measure, run_id| - type = "event_benchmarked" - tag = "writer:#{run_id}" events = - Array.new(@batch_size) do - RubyEventStore::Event.new(metadata: { event_type: type }) + @batch_size.times.map do + RubyEventStore::Event.new( + metadata: { + event_type: "event_benchmarked", + }, + ) end measure.call do - @event_store.append(events, stream_name: tag, expected_version: :any) + @event_store.append( + events, + stream_name: "writer:#{run_id}", + expected_version: :any, + ) end end end diff --git a/lib/benchmark/res_concurrent_append_conflicting_streams.rb b/lib/benchmark/res_concurrent_append_conflicting_streams.rb index bcda353..a643b9b 100644 --- a/lib/benchmark/res_concurrent_append_conflicting_streams.rb +++ b/lib/benchmark/res_concurrent_append_conflicting_streams.rb @@ -1,8 +1,5 @@ # frozen_string_literal: true -require "active_record" -require "rails_event_store" - module En57 module Benchmark Scenario.define do @@ -12,32 +9,33 @@ module En57 batch_size 100 setup do |database_url| + require "active_record" + require "rails_event_store" + ActiveRecord::Base.establish_connection(database_url) @event_store = RailsEventStore::JSONClient.new end call do |measure, run_id| - type = "event_benchmarked" - stream_name = "writer:#{run_id}" concurrently do |_writer_id, barrier| + stream_name = "writer:#{run_id}" events = - Array.new(@batch_size) do - RubyEventStore::Event.new(metadata: { event_type: type }) + @batch_size.times.map do + RubyEventStore::Event.new( + metadata: { + event_type: "event_benchmarked", + }, + ) end - position = -1 - + expected_version = -1 barrier.wait measure.call do begin - @event_store.append( - events, - stream_name: stream_name, - expected_version: position, - ) + @event_store.append(events, stream_name:, expected_version:) rescue RubyEventStore::WrongExpectedEventVersion record_retry - position = @event_store.read.stream(stream_name).count - 1 + expected_version = @event_store.read.stream(stream_name).count - 1 retry end end diff --git a/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb b/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb index 20195e2..137f34a 100644 --- a/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb +++ b/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb @@ -1,8 +1,5 @@ # frozen_string_literal: true -require "active_record" -require "rails_event_store" - module En57 module Benchmark Scenario.define do @@ -12,26 +9,29 @@ module En57 batch_size 100 setup do |database_url| + require "active_record" + require "rails_event_store" + ActiveRecord::Base.establish_connection(database_url) @event_store = RailsEventStore::JSONClient.new end call do |measure, _run_id| - type = "event_benchmarked" concurrently do |writer_id, barrier| - tag = "writer:#{writer_id}" - @event_store.read.stream(tag).of_type(type) events = - Array.new(@batch_size) do - RubyEventStore::Event.new(metadata: { event_type: type }) + @batch_size.times.map do + RubyEventStore::Event.new( + metadata: { + event_type: "event_benchmarked", + }, + ) end - barrier.wait measure.call do @event_store.append( events, - stream_name: tag, + stream_name: "writer:#{writer_id}", expected_version: :none, ) end