Andrey Yarovoy created HDDS-16846:
-------------------------------------
Summary: 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
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]