[ 
https://issues.apache.org/jira/browse/HDDS-16846?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Andrey Yarovoy reassigned HDDS-16846:
-------------------------------------

    Assignee: Andrey Yarovoy

> EC: Functional parity of the EC read path with HDFS DFSStripedInputStream
> -------------------------------------------------------------------------
>
>                 Key: HDDS-16846
>                 URL: https://issues.apache.org/jira/browse/HDDS-16846
>             Project: Apache Ozone
>          Issue Type: Epic
>          Components: EC Client
>            Reporter: Andrey Yarovoy
>            Assignee: Andrey Yarovoy
>            Priority: Major
>
> h2. Background
> Ozone EC keys are increasingly read by the same engines that read HDFS EC 
> files
> (Impala, Hive, Spark, Trino), through \{{ofs://}} / \{{o3fs://}}. Those 
> engines are
> tuned against the read semantics of HDFS \{{DFSStripedInputStream}}. This 
> umbrella
> compares the two client read paths at the code level and tracks the functional
> gaps on the Ozone side.
> References: HDFS branch-3.4 \{{DFSStripedInputStream}}, \{{StripeReader}},
> {\{StatefulStripeReader}}, \{{PositionStripeReader}}, \{{StripedBlockUtil}},
> {\{DFSInputStream}}. Ozone master: \{{ECBlockInputStreamProxy}},
> {\{ECBlockInputStream}}, \{{ECBlockReconstructedInputStream}},
> {\{ECBlockReconstructedStripeInputStream}}, \{{MultipartInputStream}},
> {\{OzoneFSInputStream}}.
>  
> h2. How HDFS reads an EC block group
> * *Healthy read is parallel per stripe.* A stateful read divides the requested
> range into \{{AlignedStripe}}s (\{{StripedBlockUtil.divideOneStripe}}) and
> {\{StripeReader.readStripe}} submits one read per needed data chunk to an
> {\{ExecutorCompletionService}} on the client-wide striped-reads pool
> (\{{dfs.client.read.striped.threadpool.size}}). The cells of a stripe are
> fetched concurrently into \{{curStripeBuf}}; a seek that lands inside the
> buffered stripe does not discard it.
> * *Degraded read is inline and per stripe.* When a chunk read fails (I/O 
> error,
> checksum error, missing location), the reader for that index is marked
> {\{shouldSkip}}, \{{prepareParityChunk}} schedules only as many parity reads 
> as
> there are missing chunks, the chunks already fetched for this stripe are kept,
> futures are consumed in *completion* order, surplus futures are cancelled once
> {\{dataBlkNum}} chunks are available, and only the missing cells are decoded
> (\{{decodeAndFillBuffer}}). Later stripes skip the bad index directly; healthy
> indexes keep being read as data.
> * *Positional read is stateless and range-exact.* \{{pread}} is not 
> synchronized
> on the stream, uses its own reader set (\{{preaderInfos}}), and
> {\{divideByteRangeIntoStripes}} reads (and, if degraded, decodes) only the 
> byte
> range requested, not whole cells.
> * *Client-detected corruption is reported.* \{{readToBuffer}} records checksum
> failures in \{{CorruptedBlocks}}; \{{readWithStrategy}} / \{{pread}} call
> {\{reportCheckSumFailure}} in \{{finally}}, which reports the replica to the
> NameNode so it is re-replicated / reconstructed.
> * *Failed datanodes are remembered across blocks of the stream* (local dead
> nodes; optionally the shared \{{DeadNodeDetector}}), and "could not obtain
> block" warnings are de-duplicated per node (\{{warnedNodes}}).
> * *Read statistics are EC-aware.* \{{ReadStatistics}} accumulates EC decoding
> time (\{{addErasureCodingDecodingTime}}), and \{{FileSystem.Statistics}} gets
> {\{bytesReadErasureCoded}} via \{{updateFileSystemECReadStats}}.
> h2. How Ozone reads an EC block group today
> * \{{ECBlockInputStreamProxy}} starts with \{{ECBlockInputStream}} when 
> enough data
> locations exist, otherwise with \{{ECBlockReconstructedInputStream}}.
> * \{{ECBlockInputStream.readWithStrategy}} reads *one cell at a time, 
> serially*:
> each loop iteration selects the current data index, opens its internal
> STANDALONE block stream, does one synchronous read bounded by the cell, then
> advances to the next index. No executor is involved on the healthy path.
> * On a read failure that spare locations cannot cover, \{{ECBlockInputStream}}
> throws \{{BadDataLocationException}}; the proxy discards the plain stream,
> switches the *remainder of the block* to the reconstruction reader
> (\{{failoverToReconstructionRead}}), rewinds the caller buffer and re-reads.
> * \{{ECBlockReconstructedStripeInputStream}} reads in parallel, but only whole
> stripes: it waits on futures in submit order (\{{loadDataBuffersFromStream}}),
> and on any error during a stripe it re-seeks and re-reads the whole stripe
> (\{{read(ByteBuffer[])}} retry loop). Its \{{seek}} must be stripe-aligned, so
> {\{ECBlockReconstructedInputStream.seek}} reads a full stripe to serve any
> offset (\{{readAndSeekStripe}}).
> * \{{MultipartInputStream.statelessPositionedReadSupported}} only accepts 
> parts
> that are all \{{BlockInputStream}} or all \{{StreamBlockInputStream}}, so EC
> keys never take the stateless pread path; \{{OzoneFSInputStream}} falls back
> to the synchronized seek-read-restore path for them.
> * Client-detected checksum mismatches cause a retry on another replica /
> reconstruction, but are not reported to SCM or the datanode; the bad replica
> is found only by the datanode container scanner.
> h2. Gap table
> ||#||Area||HDFS (branch-3.4)||Ozone (master)||
> |1|Healthy-path cell reads|Parallel per stripe, \{{StripeReader.readStripe}} 
> on the striped-reads pool|Serial, one cell per iteration, 
> \{{ECBlockInputStream.readWithStrategy}}|
> |2|Stripe buffering on the healthy path|\{{curStripeBuf}}; seek within the 
> stripe keeps it|No stripe buffer on the healthy path (the reconstruction 
> reader has one)|
> |3|Degraded read granularity|Inline, per stripe: keep fetched cells, read 
> only \{{missingChunksNum}} parity, decode only missing cells, 
> completion-order futures, cancel surplus|Whole-block failover to 
> reconstruction, rewind and re-read; any error inside a stripe re-reads the 
> whole stripe; submit-order futures|
> |4|Degraded small / unaligned reads|Range-exact 
> (\{{divideByteRangeIntoStripes}}); decodes only the requested span|Full 
> stripe read and decode for any offset (\{{readAndSeekStripe}})|
> |5|Positional read on EC|Unsynchronized, separate reader set, safe for 
> concurrent preads|Excluded from stateless pread; synchronized 
> seek-read-restore|
> |6|Corrupt replica reporting|\{{reportCheckSumFailure}} -> NameNode|None from 
> the client|
> |7|Failed datanodes across blocks of a stream|Local dead nodes (+ optional 
> \{{DeadNodeDetector}}); de-duplicated warnings|Per block stream only 
> (\{{failedLocations}}); the next block of the key starts with no knowledge of 
> nodes that already failed|
> |8|EC read statistics|\{{bytesReadErasureCoded}}, EC decoding time in 
> \{{ReadStatistics}}|Only \{{bytesRead}}; reconstruction counters are 
> client-global in \{{XceiverClientMetrics}}|
> |9|Read pool scope|One striped-reads pool for healthy and degraded reads|Pool 
> exists only for reconstruction 
> (\{{ozone.client.ec.reconstruct.stripe.read.pool.limit}})|
> h2. Already at parity (not in scope)
> * \{{ByteBuffer}} read and \{{ByteBuffer}} pread at the FileSystem level
> (\{{ByteBufferReadable}}, \{{ByteBufferPositionedReadable}}).
> * \{{StreamCapabilities}} on the FS input stream: 
> \{{CapableOzoneFSInputStream}}
> (Hadoop 3 filesystems) reports \{{READBYTEBUFFER}}, \{{UNBUFFER}},
> {\{PREADBYTEBUFFER}}.
> * \{{unbuffer}} releasing per-block resources; decoder and pooled 
> reconstruction
> buffers released on close.
> * Location refresh on read failure (\{{refreshFunction}}) and retry on spare
> replicas (\{{HDDS-7917}}).
> * Zero-copy / \{{HasEnhancedByteBufferAccess}}: unsupported in both.
> * Hedged reads: not supported for striped reads in HDFS either.
> h2. Intentional differences (keep)
> * Parity selection: HDFS takes the first usable parity index; Ozone shuffles 
> the
> candidates (\{{ECBlockReconstructedStripeInputStream}}) to spread
> reconstruction load. No change proposed.
> * Rejection policy: both striped pools run the task in the caller thread when
> saturated.
> h2. Out of scope
> * Speculative "read \{{dataBlkNum}} + 1 chunks" for full stripes (an open TODO
> in HDFS \{{StripeReader}}, not implemented there).
> * Zero-copy reads, hedged reads.
> * EC write path, offline reconstruction on the datanode.
> h2. Related
> * HDDS-5952 (resolved): parallel reads inside 
> \{{ECBlockReconstructedStripeInputStream}}.
> * HDDS-7917 (resolved): try spare replicas on error.
> * HDDS-12951 (resolved): metric/log on fallback to reconstruction reads.
> * HDDS-15424 (resolved): concurrent positional read.
> * HDDS-16396 (open): adopt PositionedReadable APIs and concurrent pread test.
> * HDDS-14761 (open): EC replica read fails once with GetBlock error after 
> offline reconstruction.



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

---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to