diff --git a/cpp/src/parquet/CMakeLists.txt b/cpp/src/parquet/CMakeLists.txt index 606dcdc0a9e..59ac9526525 100644 --- a/cpp/src/parquet/CMakeLists.txt +++ b/cpp/src/parquet/CMakeLists.txt @@ -412,6 +412,14 @@ add_parquet_test(arrow-reader-writer-test arrow/arrow_statistics_test.cc arrow/variant_test.cc) +if(ARROW_WITH_OPENTELEMETRY AND NOT ARROW_USE_ASAN) + add_parquet_test(arrow-reader-writer-tracing-test + SOURCES + arrow/arrow_reader_writer_tracing_test.cc + EXTRA_LINK_LIBS + ${ARROW_OPENTELEMETRY_LIBS}) +endif() + add_parquet_test(arrow-index-test SOURCES arrow/index_test.cc) add_parquet_test(arrow-internals-test SOURCES arrow/path_internal_test.cc diff --git a/cpp/src/parquet/arrow/arrow_reader_writer_tracing_test.cc b/cpp/src/parquet/arrow/arrow_reader_writer_tracing_test.cc new file mode 100644 index 00000000000..f9a1632aec5 --- /dev/null +++ b/cpp/src/parquet/arrow/arrow_reader_writer_tracing_test.cc @@ -0,0 +1,287 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +#include +#include +#include +#include +#include +#include +#include +#include + +#include + +#include +#include +#include +#include +#include +#include +#include +#include + +#include "arrow/api.h" +#include "arrow/io/api.h" +#include "arrow/testing/gtest_util.h" +#include "arrow/testing/util.h" +#include "parquet/arrow/reader.h" +#include "parquet/arrow/writer.h" +#include "parquet/file_reader.h" + +namespace parquet { +namespace arrow { + +namespace { + +namespace otel = opentelemetry; +namespace sdktrace = opentelemetry::sdk::trace; + +struct CapturedSpan { + std::string name; + std::unordered_map attributes; +}; + +class CapturedSpanStorage { + public: + void Clear() { + std::lock_guard lock(mutex_); + spans_.clear(); + } + + void Append(const sdktrace::SpanData& span) { + std::lock_guard lock(mutex_); + auto name = span.GetName(); + spans_.push_back({std::string(name.data(), name.size()), span.GetAttributes()}); + } + + std::vector ReadColumnSpans() const { + std::lock_guard lock(mutex_); + std::vector result; + for (const auto& span : spans_) { + if (span.name == "parquet::arrow::read_column") { + result.push_back(span); + } + } + return result; + } + + private: + mutable std::mutex mutex_; + std::vector spans_; +}; + +class CapturingSpanExporter : public sdktrace::SpanExporter { + public: + explicit CapturingSpanExporter(std::shared_ptr storage) + : storage_(std::move(storage)) {} + + std::unique_ptr MakeRecordable() noexcept override { + return std::make_unique(); + } + + otel::sdk::common::ExportResult Export( + const otel::nostd::span>& spans) noexcept + override { + try { + for (const auto& recordable : spans) { + auto* span = dynamic_cast(recordable.get()); + if (span == nullptr) { + return otel::sdk::common::ExportResult::kFailure; + } + storage_->Append(*span); + } + return otel::sdk::common::ExportResult::kSuccess; + } catch (...) { + return otel::sdk::common::ExportResult::kFailure; + } + } + + bool ForceFlush(std::chrono::microseconds) noexcept override { return true; } + + bool Shutdown(std::chrono::microseconds) noexcept override { return true; } + + private: + std::shared_ptr storage_; +}; + +const auto kSpanStorage = std::make_shared(); + +class OtelEnvironment : public ::testing::Environment { + public: + void SetUp() override { + auto exporter = std::make_unique(kSpanStorage); + auto processor = std::make_unique(std::move(exporter)); + auto provider = otel::nostd::shared_ptr( + new sdktrace::TracerProvider(std::move(processor))); + otel::trace::Provider::SetTracerProvider(std::move(provider)); + } +}; + +[[maybe_unused]] static ::testing::Environment* kOtelEnvironment = + ::testing::AddGlobalTestEnvironment(new OtelEnvironment); + +::arrow::Result> WriteToBuffer( + const std::shared_ptr<::arrow::Table>& table) { + ARROW_ASSIGN_OR_RAISE(auto sink, ::arrow::io::BufferOutputStream::Create()); + RETURN_NOT_OK( + WriteTable(*table, ::arrow::default_memory_pool(), sink, table->num_rows())); + return sink->Finish(); +} + +::arrow::Result> OpenReader( + const std::shared_ptr<::arrow::Table>& table) { + ARROW_ASSIGN_OR_RAISE(auto buffer, WriteToBuffer(table)); + auto parquet_reader = + ParquetFileReader::Open(std::make_shared<::arrow::io::BufferReader>(buffer)); + ARROW_ASSIGN_OR_RAISE(auto reader, FileReader::Make(::arrow::default_memory_pool(), + std::move(parquet_reader))); + reader->set_use_threads(false); + return reader; +} + +std::shared_ptr<::arrow::Table> NestedTable() { + auto table_schema = ::arrow::schema( + {::arrow::field("first", ::arrow::int32()), + ::arrow::field("group", + ::arrow::struct_({::arrow::field("left", ::arrow::int32()), + ::arrow::field("right", ::arrow::utf8())})), + ::arrow::field("last", ::arrow::int64())}); + return ::arrow::TableFromJSON(table_schema, + {R"([{"first": 1, "group": {"left": 10, "right": "a"}, + "last": 100}, + {"first": 2, "group": {"left": 20, "right": "b"}, + "last": 200}])"}); +} + +template +const T* GetAttribute(const CapturedSpan& span, const std::string& name) { + auto it = span.attributes.find(name); + if (it == span.attributes.end()) { + return nullptr; + } + return otel::nostd::get_if(&it->second); +} + +void AssertColumnAttributes(const CapturedSpan& span, int32_t field_index, + const std::string& field_name, + const std::string& physical_type) { + const auto* actual_index = GetAttribute(span, "parquet.arrow.columnindex"); + ASSERT_NE(actual_index, nullptr); + EXPECT_EQ(*actual_index, field_index); + + const auto* actual_name = GetAttribute(span, "parquet.arrow.columnname"); + ASSERT_NE(actual_name, nullptr); + EXPECT_EQ(*actual_name, field_name); + + const auto* actual_type = GetAttribute(span, "parquet.arrow.physicaltype"); + ASSERT_NE(actual_type, nullptr); + EXPECT_EQ(*actual_type, physical_type); +} + +TEST(ReadColumnTracing, FlatColumnSubset) { + auto table_schema = ::arrow::schema({::arrow::field("first", ::arrow::int32()), + ::arrow::field("target", ::arrow::utf8())}); + auto table = ::arrow::TableFromJSON( + table_schema, {R"([{"first": 1, "target": "a"}, {"first": 2, "target": "b"}])"}); + ASSERT_OK_AND_ASSIGN(auto reader, OpenReader(table)); + + kSpanStorage->Clear(); + ASSERT_OK_AND_ASSIGN(auto result, reader->ReadTable({1})); + ASSERT_EQ(result->num_columns(), 1); + ASSERT_TRUE(result->column(0)->Equals(table->column(1))); + + auto spans = kSpanStorage->ReadColumnSpans(); + ASSERT_EQ(spans.size(), 1); + AssertColumnAttributes(spans[0], 1, "target", "BYTE_ARRAY"); +} + +TEST(ReadColumnTracing, NestedColumnSubset) { + ASSERT_OK_AND_ASSIGN(auto reader, OpenReader(NestedTable())); + + kSpanStorage->Clear(); + ASSERT_OK_AND_ASSIGN(auto result, reader->ReadTable({2, 3})); + ASSERT_EQ(result->num_columns(), 2); + ASSERT_EQ(result->schema()->field(0)->name(), "group"); + ASSERT_EQ(result->schema()->field(1)->name(), "last"); + + auto spans = kSpanStorage->ReadColumnSpans(); + ASSERT_EQ(spans.size(), 2); + AssertColumnAttributes(spans[0], 1, "group", ""); + AssertColumnAttributes(spans[1], 2, "last", "INT64"); +} + +TEST(ReadColumnTracing, ReorderedColumnSubset) { + ASSERT_OK_AND_ASSIGN(auto reader, OpenReader(NestedTable())); + + kSpanStorage->Clear(); + ASSERT_OK_AND_ASSIGN(auto result, reader->ReadTable({3, 0, 2})); + ASSERT_EQ(result->num_columns(), 3); + + auto spans = kSpanStorage->ReadColumnSpans(); + ASSERT_EQ(spans.size(), 3); + AssertColumnAttributes(spans[0], 2, "last", "INT64"); + AssertColumnAttributes(spans[1], 0, "first", "INT32"); + AssertColumnAttributes(spans[2], 1, "group", ""); +} + +TEST(ReadColumnTracing, FullNestedSchema) { + ASSERT_OK_AND_ASSIGN(auto reader, OpenReader(NestedTable())); + + kSpanStorage->Clear(); + ASSERT_OK_AND_ASSIGN(auto result, reader->ReadTable()); + ASSERT_EQ(result->num_columns(), 3); + + auto spans = kSpanStorage->ReadColumnSpans(); + ASSERT_EQ(spans.size(), 3); + AssertColumnAttributes(spans[0], 0, "first", "INT32"); + AssertColumnAttributes(spans[1], 1, "group", ""); + AssertColumnAttributes(spans[2], 2, "last", "INT64"); +} + +TEST(ReadColumnTracing, DirectFileTopLevelFieldRead) { + auto table = NestedTable(); + ASSERT_OK_AND_ASSIGN(auto reader, OpenReader(table)); + + kSpanStorage->Clear(); + std::shared_ptr<::arrow::ChunkedArray> result; + ASSERT_OK(reader->ReadColumn(2, &result)); + ASSERT_TRUE(result->Equals(table->column(2))); + + auto spans = kSpanStorage->ReadColumnSpans(); + ASSERT_EQ(spans.size(), 1); + AssertColumnAttributes(spans[0], 2, "last", "INT64"); +} + +TEST(ReadColumnTracing, DirectRowGroupTopLevelFieldRead) { + auto table = NestedTable(); + ASSERT_OK_AND_ASSIGN(auto reader, OpenReader(table)); + + kSpanStorage->Clear(); + std::shared_ptr<::arrow::ChunkedArray> result; + ASSERT_OK(reader->RowGroup(0)->Column(1)->Read(&result)); + ASSERT_TRUE(result->Equals(table->column(1))); + + auto spans = kSpanStorage->ReadColumnSpans(); + ASSERT_EQ(spans.size(), 1); + AssertColumnAttributes(spans[0], 1, "group", ""); +} + +} // namespace + +} // namespace arrow +} // namespace parquet diff --git a/cpp/src/parquet/arrow/reader.cc b/cpp/src/parquet/arrow/reader.cc index bb8bacc7751..fe24c77bbf1 100644 --- a/cpp/src/parquet/arrow/reader.cc +++ b/cpp/src/parquet/arrow/reader.cc @@ -233,7 +233,8 @@ class FileReaderImpl : public FileReader { Status GetFieldReaders(const std::vector& column_indices, const std::vector& row_groups, std::vector>* out, - std::shared_ptr<::arrow::Schema>* out_schema) { + std::shared_ptr<::arrow::Schema>* out_schema, + std::vector* out_field_indices) { // We only need to read schema fields which have columns indicated // in the indices vector ARROW_ASSIGN_OR_RAISE(std::vector field_indices, @@ -253,6 +254,7 @@ class FileReaderImpl : public FileReader { } *out_schema = ::arrow::schema(std::move(out_fields), manifest_.schema_metadata); + *out_field_indices = std::move(field_indices); return Status::OK(); } @@ -268,8 +270,8 @@ class FileReaderImpl : public FileReader { reader_->metadata()->key_value_metadata(), out); } - Status ReadColumn(int i, const std::vector& row_groups, ColumnReader* reader, - std::shared_ptr* out) { + Status ReadColumn(int field_index, const std::vector& row_groups, + ColumnReader* reader, std::shared_ptr* out) { BEGIN_PARQUET_CATCH_EXCEPTIONS // NextBatch()'s size is a number of records (rows), not leaf values, so use the // row group's own row count directly rather than some column's num_values(). @@ -278,12 +280,18 @@ class FileReaderImpl : public FileReader { records_to_read += reader_->metadata()->RowGroup(row_group)->num_rows(); } #ifdef ARROW_WITH_OPENTELEMETRY - std::string column_name = reader_->metadata()->schema()->Column(i)->name(); - std::string phys_type = - TypeToString(reader_->metadata()->schema()->Column(i)->physical_type()); + const auto& schema_field = manifest_.schema_fields[field_index]; + const std::string& column_name = schema_field.field->name(); + std::string phys_type; + if (schema_field.is_leaf()) { + phys_type = TypeToString(reader_->metadata() + ->schema() + ->Column(schema_field.column_index) + ->physical_type()); + } ::arrow::util::tracing::Span span; START_SPAN(span, "parquet::arrow::read_column", - {{"parquet.arrow.columnindex", i}, + {{"parquet.arrow.columnindex", field_index}, {"parquet.arrow.columnname", column_name}, {"parquet.arrow.physicaltype", phys_type}, {"parquet.arrow.records_to_read", records_to_read}}); @@ -292,15 +300,16 @@ class FileReaderImpl : public FileReader { END_PARQUET_CATCH_EXCEPTIONS } - Status ReadColumn(int i, const std::vector& row_groups, + Status ReadColumn(int field_index, const std::vector& row_groups, std::shared_ptr* out) { std::unique_ptr flat_column_reader; - RETURN_NOT_OK(GetColumn(i, SomeRowGroupsFactory(row_groups), &flat_column_reader)); - return ReadColumn(i, row_groups, flat_column_reader.get(), out); + RETURN_NOT_OK( + GetColumn(field_index, SomeRowGroupsFactory(row_groups), &flat_column_reader)); + return ReadColumn(field_index, row_groups, flat_column_reader.get(), out); } - Status ReadColumn(int i, std::shared_ptr* out) override { - return ReadColumn(i, Iota(reader_->metadata()->num_row_groups()), out); + Status ReadColumn(int field_index, std::shared_ptr* out) override { + return ReadColumn(field_index, Iota(reader_->metadata()->num_row_groups()), out); } Result> ReadTable() override { @@ -1113,7 +1122,9 @@ Result> FileReaderImpl::GetRecordBatchReader( std::vector> readers; std::shared_ptr<::arrow::Schema> batch_schema; - RETURN_NOT_OK(GetFieldReaders(column_indices, row_groups, &readers, &batch_schema)); + std::vector field_indices; + RETURN_NOT_OK(GetFieldReaders(column_indices, row_groups, &readers, &batch_schema, + &field_indices)); if (readers.empty()) { // Just generate all batches right now; they're cheap since they have no columns. @@ -1373,15 +1384,18 @@ Future> FileReaderImpl::DecodeRowGroups( // in a sync context too so use `this` over `self` std::vector> readers; std::shared_ptr<::arrow::Schema> result_schema; - RETURN_NOT_OK(GetFieldReaders(column_indices, row_groups, &readers, &result_schema)); + std::vector field_indices; + RETURN_NOT_OK(GetFieldReaders(column_indices, row_groups, &readers, &result_schema, + &field_indices)); // OptionalParallelForAsync requires an executor if (!cpu_executor) cpu_executor = ::arrow::internal::GetCpuThreadPool(); - auto read_column = [row_groups, self, this](size_t i, - std::shared_ptr reader) + auto read_column = [field_indices, row_groups, self, this]( + size_t reader_index, std::shared_ptr reader) -> ::arrow::Result> { std::shared_ptr<::arrow::ChunkedArray> column; - RETURN_NOT_OK(ReadColumn(static_cast(i), row_groups, reader.get(), &column)); + RETURN_NOT_OK( + ReadColumn(field_indices[reader_index], row_groups, reader.get(), &column)); return column; }; auto make_table = [result_schema, row_groups, self, diff --git a/cpp/src/parquet/arrow/reader.h b/cpp/src/parquet/arrow/reader.h index 26269b32b93..0d60c09c10c 100644 --- a/cpp/src/parquet/arrow/reader.h +++ b/cpp/src/parquet/arrow/reader.h @@ -149,7 +149,7 @@ class PARQUET_EXPORT FileReader { /// 2 foo3 /// /// i=0 will read the entire foo struct, i=1 the foo2 primitive column etc - virtual ::arrow::Status ReadColumn(int i, + virtual ::arrow::Status ReadColumn(int field_index, std::shared_ptr<::arrow::ChunkedArray>* out) = 0; /// \brief Return a RecordBatchReader of all row groups and columns.