Skip to content

Commit 4d8d406

Browse files
committed
Fix buffer stage size accounting race
Signed-off-by: simonyang08 <ppt5928@gmail.com>
1 parent a84933f commit 4d8d406

2 files changed

Lines changed: 84 additions & 11 deletions

File tree

‎lib/fluent/plugin/buffer.rb‎

Lines changed: 10 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -383,6 +383,12 @@ def write(metadata_and_data, format: nil, size: nil, enqueue: false)
383383
if enqueue || first_chunk.unstaged? || chunk_size_full?(first_chunk)
384384
chunks_to_enqueue << first_chunk
385385
end
386+
if (bytesize = staged_bytesizes_by_chunk[first_chunk])
387+
# Account while the chunk lock is still held so enqueue_chunk
388+
# cannot remove the chunk and subtract its bytes first.
389+
@stage_size_metrics.add(bytesize)
390+
log.on_trace { log.trace { "chunk #{first_chunk.path} size_added: #{bytesize} new_size: #{first_chunk.bytesize}" } }
391+
end
386392
first_chunk.mon_exit
387393
rescue
388394
operated_chunks.unshift(first_chunk)
@@ -397,6 +403,10 @@ def write(metadata_and_data, format: nil, size: nil, enqueue: false)
397403
if enqueue || chunk.unstaged? || chunk_size_full?(chunk)
398404
chunks_to_enqueue << chunk
399405
end
406+
if (bytesize = staged_bytesizes_by_chunk[chunk])
407+
@stage_size_metrics.add(bytesize)
408+
log.on_trace { log.trace { "chunk #{chunk.path} size_added: #{bytesize} new_size: #{chunk.bytesize}" } }
409+
end
400410
chunk.mon_exit
401411
rescue => e
402412
chunk.rollback
@@ -407,17 +417,6 @@ def write(metadata_and_data, format: nil, size: nil, enqueue: false)
407417

408418
# All locks about chunks are released.
409419

410-
#
411-
# Now update the stage, stage_size with proper locking
412-
# FIX FOR stage_size miscomputation - https://github.com/fluent/fluentd/issues/2712
413-
#
414-
staged_bytesizes_by_chunk.each do |chunk, bytesize|
415-
chunk.synchronize do
416-
synchronize { @stage_size_metrics.add(bytesize) }
417-
log.on_trace { log.trace { "chunk #{chunk.path} size_added: #{bytesize} new_size: #{chunk.bytesize}" } }
418-
end
419-
end
420-
421420
chunks_to_enqueue.each do |c|
422421
if c.staged? && (enqueue || chunk_size_full?(c))
423422
m = c.metadata
Lines changed: 74 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,74 @@
1+
require_relative '../helper'
2+
require 'fluent/plugin/buffer'
3+
require 'fluent/plugin/buffer/memory_chunk'
4+
require 'fluent/plugin_id'
5+
require 'fluent/log'
6+
7+
module StageSizeRace
8+
class Owner < Fluent::Plugin::Base
9+
include Fluent::PluginId
10+
include Fluent::PluginLoggerMixin
11+
end
12+
13+
class Buf < Fluent::Plugin::Buffer
14+
def create_metadata(timekey = nil, tag = nil, variables = nil)
15+
Fluent::Plugin::Buffer::Metadata.new(timekey, tag, variables)
16+
end
17+
18+
def resume
19+
return {}, []
20+
end
21+
22+
def generate_chunk(metadata)
23+
Fluent::Plugin::Buffer::MemoryChunk.new(metadata)
24+
end
25+
end
26+
end
27+
28+
class StageSizeRaceTest < ::Test::Unit::TestCase
29+
test 'stage_size does not go negative while a write is pending' do
30+
b = StageSizeRace::Buf.new
31+
b.owner = StageSizeRace::Owner.new
32+
b.configure(config_element('buffer', '', { 'total_limit_size' => 1024, 'chunk_limit_size' => 4096 }))
33+
b.start
34+
35+
m = b.create_metadata
36+
b.write({ m => ['a' * 400] })
37+
chunk = b.stage[m]
38+
assert_equal 400, b.stage_size
39+
40+
reached = Queue.new
41+
resume = Queue.new
42+
armed = false
43+
44+
chunk.define_singleton_method(:mon_exit) do
45+
result = super()
46+
if armed
47+
armed = false
48+
reached << true
49+
resume.pop
50+
end
51+
result
52+
end
53+
54+
armed = true
55+
writer = Thread.new { b.write({ m => ['b' * 400] }) }
56+
reached.pop
57+
58+
b.enqueue_chunk(m)
59+
assert_equal 0, b.stage_size
60+
61+
resume << true
62+
writer.join
63+
64+
assert_equal 0, b.stage.size
65+
assert_equal 0, b.stage_size
66+
assert_equal 800, b.queue_size
67+
ensure
68+
if writer&.alive?
69+
resume << true
70+
writer.join
71+
end
72+
b&.stop
73+
end
74+
end

0 commit comments

Comments
 (0)