Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
refactor write worker for better performance
- Loading branch information
1 parent
d7ee2cf
commit d969206
Showing
13 changed files
with
134 additions
and
109 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
File renamed without changes.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,58 @@ | ||
require "./helpers/table_writer" | ||
require "../utils/wait_group" | ||
require "../storage/table/abstract_table" | ||
require "./executor/datum" | ||
|
||
class ReQL::Worker | ||
WORKER_COUNT = 128 | ||
CHANNEL_SIZE = 10 * WORKER_COUNT | ||
|
||
@inserter = Channel({wait_group: WaitGroup, table_writer: TableWriter, table: Storage::AbstractTable, row: Hash(String, Datum)}).new(CHANNEL_SIZE) | ||
@deleter = Channel({wait_group: WaitGroup, table_writer: TableWriter, table: Storage::AbstractTable, key: Datum}).new(CHANNEL_SIZE) | ||
@closer = Channel(WaitGroup).new(WORKER_COUNT) | ||
|
||
def initialize | ||
WORKER_COUNT.times do | ||
spawn worker | ||
end | ||
end | ||
|
||
def close | ||
wait_group = WaitGroup.new | ||
WORKER_COUNT.times do | ||
wait_group.add | ||
@closer.send(wait_group) | ||
end | ||
wait_group.wait | ||
|
||
@inserter.close | ||
@closer.close | ||
end | ||
|
||
def insert(wait_group, table_writer, table, row) | ||
@inserter.send({wait_group: wait_group, table_writer: table_writer, table: table, row: row}) | ||
end | ||
|
||
def delete(wait_group, table_writer, table, key) | ||
@deleter.send({wait_group: wait_group, table_writer: table_writer, table: table, key: key}) | ||
end | ||
|
||
private def worker | ||
loop do | ||
select | ||
when wait_group = @closer.receive | ||
wait_group.done | ||
return | ||
when tuple = @inserter.receive | ||
tuple[:table_writer].insert(tuple[:table], tuple[:row]) | ||
tuple[:wait_group].done | ||
when tuple = @deleter.receive | ||
tuple[:table_writer].delete(tuple[:table], tuple[:key]) | ||
tuple[:wait_group].done | ||
end | ||
end | ||
rescue error | ||
STDERR.puts error.inspect_with_backtrace | ||
worker | ||
end | ||
end |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
File renamed without changes.