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