Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

Configurable deadlock detector interval for ingester. (resubmit) #1134

Merged
merged 6 commits into from
Nov 12, 2018
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Prev Previous commit
Next Next commit
Making channels for deadlock detector only when it is enabled
Signed-off-by: Chodor Marek <[email protected]>
  • Loading branch information
Chodor Marek committed Oct 26, 2018
commit 1730b8bf1cff9077f9433941e2d7a48a908a7397
51 changes: 27 additions & 24 deletions cmd/ingester/app/consumer/deadlock_detector.go
Original file line number Diff line number Diff line change
Expand Up @@ -52,14 +52,14 @@ type partitionDeadlockDetector struct {
closePartition chan struct{}
done chan struct{}
incrementAllPartitionMsgCount func()
closed bool
disabled bool
}

type allPartitionsDeadlockDetector struct {
msgConsumed *uint64
logger *zap.Logger
done chan struct{}
closed bool
disabled bool
}

func newDeadlockDetector(metricsFactory metrics.Factory, logger *zap.Logger, interval time.Duration) deadlockDetector {
Expand All @@ -84,12 +84,10 @@ func newDeadlockDetector(metricsFactory metrics.Factory, logger *zap.Logger, int
func (s *deadlockDetector) startMonitoringForPartition(partition int32) *partitionDeadlockDetector {
var msgConsumed uint64
w := &partitionDeadlockDetector{
msgConsumed: &msgConsumed,
partition: partition,
closePartition: make(chan struct{}, 1),
done: make(chan struct{}),
logger: s.logger,
closed: false,
msgConsumed: &msgConsumed,
partition: partition,
logger: s.logger,
disabled: 0 == s.interval,

incrementAllPartitionMsgCount: func() {
s.allPartitionsDeadlockDetector.incrementMsgCount()
Expand All @@ -98,8 +96,9 @@ func (s *deadlockDetector) startMonitoringForPartition(partition int32) *partiti

if s.interval == 0 {
s.logger.Debug("Partition deadlock detector disabled")
w.closed = true
} else {
w.closePartition = make(chan struct{}, 1)
w.done = make(chan struct{})
go s.monitorForPartition(w, partition)
}

Expand Down Expand Up @@ -141,16 +140,15 @@ func (s *deadlockDetector) start() {
var msgConsumed uint64
detector := &allPartitionsDeadlockDetector{
msgConsumed: &msgConsumed,
done: make(chan struct{}),
logger: s.logger,
closed: false,
disabled: 0 == s.interval,
}

if s.interval == 0 {
s.logger.Debug("Global deadlock detector disabled")
detector.closed = true
} else {
s.logger.Debug("Starting global deadlock detector")
detector.done = make(chan struct{})
go func() {
ticker := time.NewTicker(s.interval)
defer ticker.Stop()
Expand All @@ -175,34 +173,39 @@ func (s *deadlockDetector) start() {
}

func (s *deadlockDetector) close() {
if !s.allPartitionsDeadlockDetector.closed {
s.logger.Debug("Closing all partitions deadlock detector")
s.allPartitionsDeadlockDetector.closed = true
s.allPartitionsDeadlockDetector.done <- struct{}{}
} else {
s.logger.Debug("All partitions deadlock detector already closed")
if s.allPartitionsDeadlockDetector.disabled {
return
}
s.logger.Debug("Closing all partitions deadlock detector")
s.allPartitionsDeadlockDetector.done <- struct{}{}
}

func (s *allPartitionsDeadlockDetector) incrementMsgCount() {
if s.disabled {
return
}
atomic.AddUint64(s.msgConsumed, 1)
}

func (w *partitionDeadlockDetector) closePartitionChannel() chan struct{} {
if w.disabled {
return nil
}
return w.closePartition
}

func (w *partitionDeadlockDetector) close() {
if !w.closed {
w.logger.Debug("Closing deadlock detector", zap.Int32("partition", w.partition))
w.closed = true
w.done <- struct{}{}
} else {
w.logger.Debug("Deadlock detector already closed", zap.Int32("partition", w.partition))
if w.disabled {
return
}
w.logger.Debug("Closing deadlock detector", zap.Int32("partition", w.partition))
w.done <- struct{}{}
}

func (w *partitionDeadlockDetector) incrementMsgCount() {
if w.disabled {
return
}
w.incrementAllPartitionMsgCount()
atomic.AddUint64(w.msgConsumed, 1)
}
1 change: 1 addition & 0 deletions cmd/ingester/app/consumer/deadlock_detector_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -159,6 +159,7 @@ func TestApiCompatibilityWhenDeadlockDetectorDisabled(t *testing.T) {
w := f.startMonitoringForPartition(1)

w.incrementMsgCount()
w.incrementAllPartitionMsgCount()
assert.Zero(t, len(w.closePartitionChannel()))
w.close()
Copy link
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

I am not clear what this is testing. Why calling increments should have any effect on the closePartitionChannel?

}