@@ -71,46 +71,70 @@ def add_or_update(feature_name, key)
7171 end
7272 end
7373
74+ def clear_cache
75+ uniques = @cache . clone
76+ keys_size = @keys_size
77+ @cache . clear
78+ @keys_size = 0
79+
80+ [ uniques , keys_size ]
81+ end
82+
7483 def send_bulk_data
7584 @semaphore . synchronize do
7685 return if @cache . empty?
7786
78- uniques = @cache . clone
79- keys_size = @keys_size
80- @cache . clear
81- @keys_size = 0
82-
87+ uniques , keys_size = clear_cache
8388 if keys_size <= @max_bulk_size
8489 @sender_adapter . record_uniques_key ( uniques )
8590 return
86- end
87-
88- bulks = [ ]
89- uniques . each do |unique |
90- bulks += check_keys_and_split_to_bulks ( unique )
91- end
9291
93- bulks . each do |b |
94- @sender_adapter . record_uniques_key ( b )
9592 end
93+ bulks = flatten_bulks ( uniques )
94+ bulks_to_post = group_bulks_by_max_size ( bulks )
95+ @sender_adapter . record_uniques_key ( bulks_to_post )
9696 end
9797 rescue StandardError => e
9898 @config . log_found_exception ( __method__ . to_s , e )
9999 end
100100
101- def check_keys_and_split_to_bulks ( unique )
102- unique_updated = [ ]
103- unique . each do |_ , value |
104- if value . size > @max_bulk_size
105- sub_bulks = SplitIoClient ::Utilities . split_bulk_to_send ( value , value . size / @max_bulk_size )
106- sub_bulks . each do |sub_bulk |
107- unique_updated . add ( { key : sub_bulk } )
108- end
109- break
101+ def group_bulks_by_max_size ( bulks )
102+ current_size = 0
103+ bulks_to_post = Concurrent ::Hash . new
104+ bulks . each do |bulk |
105+ key , value = bulk . first
106+ if ( value . size + current_size ) > @max_bulk_size
107+ @sender_adapter . record_uniques_key ( bulks_to_post )
108+ bulks_to_post = Concurrent ::Hash . new
109+ current_size = 0
110+ end
111+ bulks_to_post [ key ] = value
112+ current_size += value . size
113+ end
114+
115+ bulks_to_post
116+ end
117+
118+ def flatten_bulks ( uniques )
119+ bulks = [ ]
120+ uniques . each_key do |unique_key |
121+ bulks += check_keys_and_split_to_bulks ( uniques [ unique_key ] , unique_key )
122+ end
110123
124+ bulks
125+ end
126+
127+ def check_keys_and_split_to_bulks ( value , key )
128+ unique_updated = [ ]
129+ if value . size > @max_bulk_size
130+ sub_bulks = SplitIoClient ::Utilities . split_bulk_to_send ( value , @max_bulk_size )
131+ sub_bulks . each do |sub_bulk |
132+ unique_updated << { key => sub_bulk . to_set }
111133 end
112- unique_updated . add ( { key : value } )
134+ return unique_updated
135+
113136 end
137+ unique_updated << { key => value }
114138
115139 unique_updated
116140 end
0 commit comments