MCPcopy Create free account
hub / github.com/apache/arrow / InitFromBlock

Method InitFromBlock

cpp/src/arrow/csv/reader.cc:912–958  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

910 }
911
912 Future<> InitFromBlock(const DecodedBlock& block,
913 AsyncGenerator<DecodedBlock> batch_gen, int max_readahead,
914 int64_t prev_bytes_processed) {
915 if (!block.record_batch) {
916 // End of file just return null batches
917 record_batch_gen_ = MakeEmptyGenerator<std::shared_ptr<RecordBatch>>();
918 return Status::OK();
919 }
920
921 schema_ = block.record_batch->schema();
922
923 if (block.record_batch->num_rows() == 0) {
924 // Keep consuming blocks until the first non empty block is found
925 auto self = shared_from_this();
926 prev_bytes_processed += block.bytes_processed;
927 return batch_gen().Then([self, batch_gen, max_readahead,
928 prev_bytes_processed](const DecodedBlock& next_block) {
929 return self->InitFromBlock(next_block, std::move(batch_gen), max_readahead,
930 prev_bytes_processed);
931 });
932 }
933
934 AsyncGenerator<DecodedBlock> readahead_gen;
935 if (read_options_.use_threads) {
936 readahead_gen = MakeReadaheadGenerator(std::move(batch_gen), max_readahead);
937 } else {
938 readahead_gen = std::move(batch_gen);
939 }
940
941 AsyncGenerator<DecodedBlock> restarted_gen =
942 MakeGeneratorStartsWith({block}, std::move(readahead_gen));
943
944 auto bytes_decoded = bytes_decoded_;
945 auto unwrap_and_record_bytes =
946 [bytes_decoded, prev_bytes_processed](
947 const DecodedBlock& block) mutable -> Result<std::shared_ptr<RecordBatch>> {
948 bytes_decoded->fetch_add(block.bytes_processed + prev_bytes_processed);
949 prev_bytes_processed = 0;
950 return block.record_batch;
951 };
952
953 auto unwrapped =
954 MakeMappedGenerator(std::move(restarted_gen), std::move(unwrap_and_record_bytes));
955
956 record_batch_gen_ = MakeCancellable(std::move(unwrapped), io_context_.stop_token());
957 return Status::OK();
958 }
959
960 std::shared_ptr<Schema> schema_;
961 AsyncGenerator<std::shared_ptr<RecordBatch>> record_batch_gen_;

Callers 1

InitAfterFirstBufferMethod · 0.95

Calls 9

MakeReadaheadGeneratorFunction · 0.85
MakeGeneratorStartsWithFunction · 0.85
MakeMappedGeneratorFunction · 0.85
MakeCancellableFunction · 0.85
ThenMethod · 0.80
OKFunction · 0.50
batch_genFunction · 0.50
schemaMethod · 0.45
num_rowsMethod · 0.45

Tested by

no test coverage detected