Alexey Serbin created KUDU-3814:
-----------------------------------

             Summary: Explore alternative approaches of flushing data to 
persistent storage
                 Key: KUDU-3814
                 URL: https://issues.apache.org/jira/browse/KUDU-3814
             Project: Kudu
          Issue Type: Task
            Reporter: Alexey Serbin


There are several places in the Kudu's code where the data is synchronized with 
the backing storage to guarantee that it's persistent.  This is done 
at-the-spot, synchronously, so right after 
{{fdatasync()/fsync()/sync_file_range()}} there is almost certain guarantee 
(modulo device controller's issues) the data isn't lost if the Kudu server 
process or the whole machine crashes.  In an active Kudu cluster, there is a 
lot of concurrent activity that's run by Kudu itself -- I/O from other 
concurrent data ingestion streams, I/O from concurrent reads/scans, background 
compaction activity, I/O induced by Raft consensus internal book-keeping when 
processing leadership changes and initiating vote for a new leader, etc.  
Sometimes it leads to I/O saturation, and oftentimes such peak conditions are 
prone to spiraling down (the root cause behind the tendency of spiraling down 
in certain scenarios is a separate item).  If the backing storage isn't 
saturated most of the time, and I/O saturation occurs only for relatively short 
time intervals (often prolonged too much due to the tendency of spiraling down 
in certain scenarios), one might think of spreading out the I/O demand induced 
by Kudu operations over the time, so the I/O demand doesn't cross the 
saturation threshold.  For example, those FS synchronization requests might be 
a queued (one queue per backing storage device), and processed at the times 
when device I/O pressure drops.

Of course, this approach means sacrificing data consistency guarantees in case 
if Kudu server process crashes or the whole machine crashes due to a hardware 
failure.  However, since Kudu's data is usually replicated (RF=3 or more) 
across multiple nodes in the cluster, for some environments it may be feasible 
to operate this way, completely replacing the data under a crashed tablet 
server or whole node, re-replicating the data from other Kudu tablet replicas 
in the cluster.  This behavior should be configurable, of course.

It's necessary to explore this approach.

Also, it seems [KUDU-484|https://issues.apache.org/jira/browse/KUDU-484] has 
some relevance in this context.





--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to