diff --git a/c_glib/parquet-glib/arrow-file-reader.cpp b/c_glib/parquet-glib/arrow-file-reader.cpp index 86bf284d1236..2db6b313f5b4 100644 --- a/c_glib/parquet-glib/arrow-file-reader.cpp +++ b/c_glib/parquet-glib/arrow-file-reader.cpp @@ -25,6 +25,31 @@ #include +namespace { + GParquetArrowFileReader * + open_reader(std::shared_ptr source, + GArrowSeekableInputStream *source_object, + GParquetReaderProperties *properties, + GError **error, + const char *tag) + { + auto parquet_properties = properties ? *gparquet_reader_properties_get_raw(properties) + : parquet::default_reader_properties(); + parquet::arrow::FileReaderBuilder builder; + if (!garrow::check(error, builder.Open(source, parquet_properties), tag)) { + return NULL; + } + if (properties) { + builder.properties(*gparquet_reader_properties_get_arrow_raw(properties)); + } + auto result = builder.Build(); + if (!garrow::check(error, result, tag)) { + return NULL; + } + return gparquet_arrow_file_reader_new_raw(result->release(), source_object); + } +} // namespace + G_BEGIN_DECLS /** @@ -32,18 +57,190 @@ G_BEGIN_DECLS * @short_description: Arrow file reader class * @include: parquet-glib/parquet-glib.h * + * #GParquetReaderProperties is a class for configuring Parquet reads. + * * #GParquetArrowFileReader is a class for reading Apache Parquet data * from file and returns them as Apache Arrow data. */ +typedef struct GParquetReaderPropertiesPrivate_ +{ + parquet::ReaderProperties properties; + parquet::ArrowReaderProperties arrow_properties; +} GParquetReaderPropertiesPrivate; + +G_DEFINE_TYPE_WITH_PRIVATE(GParquetReaderProperties, + gparquet_reader_properties, + G_TYPE_OBJECT) + +#define GPARQUET_READER_PROPERTIES_GET_PRIVATE(object) \ + static_cast( \ + gparquet_reader_properties_get_instance_private(GPARQUET_READER_PROPERTIES(object))) + +static void +gparquet_reader_properties_finalize(GObject *object) +{ + auto priv = GPARQUET_READER_PROPERTIES_GET_PRIVATE(object); + priv->arrow_properties.~ArrowReaderProperties(); + priv->properties.~ReaderProperties(); + G_OBJECT_CLASS(gparquet_reader_properties_parent_class)->finalize(object); +} + +static void +gparquet_reader_properties_init(GParquetReaderProperties *object) +{ + auto priv = GPARQUET_READER_PROPERTIES_GET_PRIVATE(object); + new (&priv->properties) parquet::ReaderProperties(parquet::default_reader_properties()); + new (&priv->arrow_properties) + parquet::ArrowReaderProperties(parquet::default_arrow_reader_properties()); +} + +static void +gparquet_reader_properties_class_init(GParquetReaderPropertiesClass *klass) +{ + G_OBJECT_CLASS(klass)->finalize = gparquet_reader_properties_finalize; +} + +/** + * gparquet_reader_properties_new: + * + * Returns: A newly created #GParquetReaderProperties. + * + * Since: 26.0.0 + */ +GParquetReaderProperties * +gparquet_reader_properties_new(void) +{ + return GPARQUET_READER_PROPERTIES(g_object_new(GPARQUET_TYPE_READER_PROPERTIES, NULL)); +} + +/** + * gparquet_reader_properties_enable_buffered_stream: + * @properties: A #GParquetReaderProperties. + * + * Enable buffered stream reading. To avoid pre-buffering whole column chunks, + * also call gparquet_reader_properties_set_pre_buffer() with %FALSE. + * This does not impose a limit on the memory used by decoded data. + * + * Since: 26.0.0 + */ +void +gparquet_reader_properties_enable_buffered_stream(GParquetReaderProperties *properties) +{ + auto priv = GPARQUET_READER_PROPERTIES_GET_PRIVATE(properties); + priv->properties.enable_buffered_stream(); +} + +/** + * gparquet_reader_properties_disable_buffered_stream: + * @properties: A #GParquetReaderProperties. + * + * Disable buffered stream reading. This is the default. + * + * Since: 26.0.0 + */ +void +gparquet_reader_properties_disable_buffered_stream(GParquetReaderProperties *properties) +{ + auto priv = GPARQUET_READER_PROPERTIES_GET_PRIVATE(properties); + priv->properties.disable_buffered_stream(); +} + +/** + * gparquet_reader_properties_is_buffered_stream_enabled: + * @properties: A #GParquetReaderProperties. + * + * Returns: %TRUE if buffered stream reading is enabled, %FALSE otherwise. + * + * Since: 26.0.0 + */ +gboolean +gparquet_reader_properties_is_buffered_stream_enabled( + GParquetReaderProperties *properties) +{ + auto priv = GPARQUET_READER_PROPERTIES_GET_PRIVATE(properties); + return priv->properties.is_buffered_stream_enabled(); +} + +/** + * gparquet_reader_properties_set_buffer_size: + * @properties: A #GParquetReaderProperties. + * @buffer_size: The buffer size in bytes. This must be positive when buffering is + * enabled. + * + * Set the buffered stream size. This does not enable buffered stream reading. + * Reads required for a data page may exceed this size. + * + * Since: 26.0.0 + */ +void +gparquet_reader_properties_set_buffer_size(GParquetReaderProperties *properties, + gint64 buffer_size) +{ + auto priv = GPARQUET_READER_PROPERTIES_GET_PRIVATE(properties); + priv->properties.set_buffer_size(buffer_size); +} + +/** + * gparquet_reader_properties_get_buffer_size: + * @properties: A #GParquetReaderProperties. + * + * Returns: The buffered stream size in bytes. + * + * Since: 26.0.0 + */ +gint64 +gparquet_reader_properties_get_buffer_size(GParquetReaderProperties *properties) +{ + auto priv = GPARQUET_READER_PROPERTIES_GET_PRIVATE(properties); + return priv->properties.buffer_size(); +} + +/** + * gparquet_reader_properties_set_pre_buffer: + * @properties: A #GParquetReaderProperties. + * @pre_buffer: Whether to pre-buffer column chunks. + * + * Set whether to pre-buffer column chunks to coalesce reads. This is enabled + * by default to improve performance on high-latency filesystems. Set this to + * %FALSE to use buffered streams without pre-buffering whole column chunks. + * This does not enable or disable buffered stream reading. + * + * Since: 26.0.0 + */ +void +gparquet_reader_properties_set_pre_buffer(GParquetReaderProperties *properties, + gboolean pre_buffer) +{ + auto priv = GPARQUET_READER_PROPERTIES_GET_PRIVATE(properties); + priv->arrow_properties.set_pre_buffer(pre_buffer); +} + +/** + * gparquet_reader_properties_get_pre_buffer: + * @properties: A #GParquetReaderProperties. + * + * Returns: %TRUE if pre-buffering is enabled, %FALSE otherwise. + * + * Since: 26.0.0 + */ +gboolean +gparquet_reader_properties_get_pre_buffer(GParquetReaderProperties *properties) +{ + auto priv = GPARQUET_READER_PROPERTIES_GET_PRIVATE(properties); + return priv->arrow_properties.pre_buffer(); +} + typedef struct GParquetArrowFileReaderPrivate_ { parquet::arrow::FileReader *arrow_file_reader; + GArrowSeekableInputStream *source; } GParquetArrowFileReaderPrivate; enum { PROP_0, - PROP_ARROW_FILE_READER + PROP_ARROW_FILE_READER, + PROP_SOURCE }; G_DEFINE_TYPE_WITH_PRIVATE(GParquetArrowFileReader, @@ -54,6 +251,14 @@ G_DEFINE_TYPE_WITH_PRIVATE(GParquetArrowFileReader, static_cast( \ gparquet_arrow_file_reader_get_instance_private(GPARQUET_ARROW_FILE_READER(obj))) +static void +gparquet_arrow_file_reader_dispose(GObject *object) +{ + auto priv = GPARQUET_ARROW_FILE_READER_GET_PRIVATE(object); + g_clear_object(&priv->source); + G_OBJECT_CLASS(gparquet_arrow_file_reader_parent_class)->dispose(object); +} + static void gparquet_arrow_file_reader_finalize(GObject *object) { @@ -73,6 +278,9 @@ gparquet_arrow_file_reader_set_property(GObject *object, auto priv = GPARQUET_ARROW_FILE_READER_GET_PRIVATE(object); switch (prop_id) { + case PROP_SOURCE: + priv->source = GARROW_SEEKABLE_INPUT_STREAM(g_value_dup_object(value)); + break; case PROP_ARROW_FILE_READER: priv->arrow_file_reader = static_cast(g_value_get_pointer(value)); @@ -89,7 +297,12 @@ gparquet_arrow_file_reader_get_property(GObject *object, GValue *value, GParamSpec *pspec) { + auto priv = GPARQUET_ARROW_FILE_READER_GET_PRIVATE(object); + switch (prop_id) { + case PROP_SOURCE: + g_value_set_object(value, priv->source); + break; default: G_OBJECT_WARN_INVALID_PROPERTY_ID(object, prop_id, pspec); break; @@ -108,6 +321,7 @@ gparquet_arrow_file_reader_class_init(GParquetArrowFileReaderClass *klass) auto gobject_class = G_OBJECT_CLASS(klass); + gobject_class->dispose = gparquet_arrow_file_reader_dispose; gobject_class->finalize = gparquet_arrow_file_reader_finalize; gobject_class->set_property = gparquet_arrow_file_reader_set_property; gobject_class->get_property = gparquet_arrow_file_reader_get_property; @@ -118,6 +332,21 @@ gparquet_arrow_file_reader_class_init(GParquetArrowFileReaderClass *klass) "The raw parquet::arrow::FileReader *", static_cast(G_PARAM_WRITABLE | G_PARAM_CONSTRUCT_ONLY)); g_object_class_install_property(gobject_class, PROP_ARROW_FILE_READER, spec); + + /** + * GParquetArrowFileReader:source: + * + * The source stream for this reader. + * + * Since: 26.0.0 + */ + spec = g_param_spec_object( + "source", + "Source", + "The source stream for this reader", + GARROW_TYPE_SEEKABLE_INPUT_STREAM, + static_cast(G_PARAM_READWRITE | G_PARAM_CONSTRUCT_ONLY)); + g_object_class_install_property(gobject_class, PROP_SOURCE, spec); } /** @@ -132,18 +361,11 @@ gparquet_arrow_file_reader_class_init(GParquetArrowFileReaderClass *klass) GParquetArrowFileReader * gparquet_arrow_file_reader_new_arrow(GArrowSeekableInputStream *source, GError **error) { - auto arrow_random_access_file = garrow_seekable_input_stream_get_raw(source); - auto arrow_memory_pool = arrow::default_memory_pool(); - auto parquet_arrow_file_reader_result = - parquet::arrow::OpenFile(arrow_random_access_file, arrow_memory_pool); - if (garrow::check(error, - parquet_arrow_file_reader_result, - "[parquet][arrow][file-reader][new-arrow]")) { - return gparquet_arrow_file_reader_new_raw( - parquet_arrow_file_reader_result->release()); - } else { - return NULL; - } + return open_reader(garrow_seekable_input_stream_get_raw(source), + source, + nullptr, + error, + "[parquet][arrow][file-reader][new-arrow]"); } /** @@ -158,27 +380,64 @@ gparquet_arrow_file_reader_new_arrow(GArrowSeekableInputStream *source, GError * GParquetArrowFileReader * gparquet_arrow_file_reader_new_path(const gchar *path, GError **error) { - auto arrow_memory_mapped_file = - arrow::io::MemoryMappedFile::Open(path, arrow::io::FileMode::READ); - if (!garrow::check(error, - arrow_memory_mapped_file, - "[parquet][arrow][file-reader][new-path]")) { + const char *tag = "[parquet][arrow][file-reader][new-path]"; + auto source = arrow::io::MemoryMappedFile::Open(path, arrow::io::FileMode::READ); + if (!garrow::check(error, source, tag)) { return NULL; } + return open_reader(*source, nullptr, nullptr, error, tag); +} - std::shared_ptr arrow_random_access_file = - arrow_memory_mapped_file.ValueOrDie(); - auto arrow_memory_pool = arrow::default_memory_pool(); - auto parquet_arrow_file_reader_result = - parquet::arrow::OpenFile(arrow_random_access_file, arrow_memory_pool); - if (garrow::check(error, - parquet_arrow_file_reader_result, - "[parquet][arrow][file-reader][new-path]")) { - return gparquet_arrow_file_reader_new_raw( - parquet_arrow_file_reader_result->release()); - } else { +/** + * gparquet_arrow_file_reader_new_arrow_full: + * @source: Arrow source to be read. + * @properties: (nullable): Reader properties or %NULL for the defaults. + * @error: (nullable): Return location for a #GError or %NULL. + * + * The reader copies @properties at construction. Later changes to @properties + * do not affect the reader. The native source is retained by the reader. + * + * Returns: (nullable): A newly created #GParquetArrowFileReader. + * + * Since: 26.0.0 + */ +GParquetArrowFileReader * +gparquet_arrow_file_reader_new_arrow_full(GArrowSeekableInputStream *source, + GParquetReaderProperties *properties, + GError **error) +{ + return open_reader(garrow_seekable_input_stream_get_raw(source), + source, + properties, + error, + "[parquet][arrow][file-reader][new-arrow-full]"); +} + +/** + * gparquet_arrow_file_reader_new_path_full: + * @path: Path to be read. + * @properties: (nullable): Reader properties or %NULL for the defaults. + * @error: (nullable): Return location for a #GError or %NULL. + * + * The reader copies @properties at construction. Later changes to @properties + * do not affect the reader. The file is memory mapped, as with + * gparquet_arrow_file_reader_new_path(). + * + * Returns: (nullable): A newly created #GParquetArrowFileReader. + * + * Since: 26.0.0 + */ +GParquetArrowFileReader * +gparquet_arrow_file_reader_new_path_full(const gchar *path, + GParquetReaderProperties *properties, + GError **error) +{ + const char *tag = "[parquet][arrow][file-reader][new-path-full]"; + auto source = arrow::io::MemoryMappedFile::Open(path, arrow::io::FileMode::READ); + if (!garrow::check(error, source, tag)) { return NULL; } + return open_reader(*source, nullptr, properties, error, tag); } /** @@ -410,12 +669,15 @@ gparquet_arrow_file_reader_get_metadata(GParquetArrowFileReader *reader) G_END_DECLS GParquetArrowFileReader * -gparquet_arrow_file_reader_new_raw(parquet::arrow::FileReader *parquet_arrow_file_reader) +gparquet_arrow_file_reader_new_raw(parquet::arrow::FileReader *parquet_arrow_file_reader, + GArrowSeekableInputStream *source) { auto arrow_file_reader = GPARQUET_ARROW_FILE_READER(g_object_new(GPARQUET_TYPE_ARROW_FILE_READER, "arrow-file-reader", parquet_arrow_file_reader, + "source", + source, NULL)); return arrow_file_reader; } @@ -426,3 +688,17 @@ gparquet_arrow_file_reader_get_raw(GParquetArrowFileReader *arrow_file_reader) auto priv = GPARQUET_ARROW_FILE_READER_GET_PRIVATE(arrow_file_reader); return priv->arrow_file_reader; } + +parquet::ReaderProperties * +gparquet_reader_properties_get_raw(GParquetReaderProperties *properties) +{ + auto priv = GPARQUET_READER_PROPERTIES_GET_PRIVATE(properties); + return &(priv->properties); +} + +parquet::ArrowReaderProperties * +gparquet_reader_properties_get_arrow_raw(GParquetReaderProperties *properties) +{ + auto priv = GPARQUET_READER_PROPERTIES_GET_PRIVATE(properties); + return &(priv->arrow_properties); +} diff --git a/c_glib/parquet-glib/arrow-file-reader.h b/c_glib/parquet-glib/arrow-file-reader.h index 04df24b87d31..93fdba586959 100644 --- a/c_glib/parquet-glib/arrow-file-reader.h +++ b/c_glib/parquet-glib/arrow-file-reader.h @@ -23,6 +23,49 @@ G_BEGIN_DECLS +#define GPARQUET_TYPE_READER_PROPERTIES (gparquet_reader_properties_get_type()) +GPARQUET_AVAILABLE_IN_26_0 +G_DECLARE_DERIVABLE_TYPE(GParquetReaderProperties, + gparquet_reader_properties, + GPARQUET, + READER_PROPERTIES, + GObject) +struct _GParquetReaderPropertiesClass +{ + GObjectClass parent_class; +}; + +GPARQUET_AVAILABLE_IN_26_0 +GParquetReaderProperties * +gparquet_reader_properties_new(void); + +GPARQUET_AVAILABLE_IN_26_0 +void +gparquet_reader_properties_enable_buffered_stream(GParquetReaderProperties *properties); +GPARQUET_AVAILABLE_IN_26_0 +void +gparquet_reader_properties_disable_buffered_stream(GParquetReaderProperties *properties); +GPARQUET_AVAILABLE_IN_26_0 +gboolean +gparquet_reader_properties_is_buffered_stream_enabled( + GParquetReaderProperties *properties); + +GPARQUET_AVAILABLE_IN_26_0 +void +gparquet_reader_properties_set_buffer_size(GParquetReaderProperties *properties, + gint64 buffer_size); +GPARQUET_AVAILABLE_IN_26_0 +gint64 +gparquet_reader_properties_get_buffer_size(GParquetReaderProperties *properties); + +GPARQUET_AVAILABLE_IN_26_0 +void +gparquet_reader_properties_set_pre_buffer(GParquetReaderProperties *properties, + gboolean pre_buffer); +GPARQUET_AVAILABLE_IN_26_0 +gboolean +gparquet_reader_properties_get_pre_buffer(GParquetReaderProperties *properties); + #define GPARQUET_TYPE_ARROW_FILE_READER (gparquet_arrow_file_reader_get_type()) GPARQUET_AVAILABLE_IN_0_11 G_DECLARE_DERIVABLE_TYPE(GParquetArrowFileReader, @@ -43,6 +86,18 @@ GPARQUET_AVAILABLE_IN_0_11 GParquetArrowFileReader * gparquet_arrow_file_reader_new_path(const gchar *path, GError **error); +GPARQUET_AVAILABLE_IN_26_0 +GParquetArrowFileReader * +gparquet_arrow_file_reader_new_arrow_full(GArrowSeekableInputStream *source, + GParquetReaderProperties *properties, + GError **error); + +GPARQUET_AVAILABLE_IN_26_0 +GParquetArrowFileReader * +gparquet_arrow_file_reader_new_path_full(const gchar *path, + GParquetReaderProperties *properties, + GError **error); + GPARQUET_AVAILABLE_IN_23_0 void gparquet_arrow_file_reader_close(GParquetArrowFileReader *reader); diff --git a/c_glib/parquet-glib/arrow-file-reader.hpp b/c_glib/parquet-glib/arrow-file-reader.hpp index 172dcccb0d45..451ccbb7c829 100644 --- a/c_glib/parquet-glib/arrow-file-reader.hpp +++ b/c_glib/parquet-glib/arrow-file-reader.hpp @@ -24,6 +24,11 @@ #include GParquetArrowFileReader * -gparquet_arrow_file_reader_new_raw(parquet::arrow::FileReader *parquet_arrow_file_reader); +gparquet_arrow_file_reader_new_raw(parquet::arrow::FileReader *parquet_arrow_file_reader, + GArrowSeekableInputStream *source = nullptr); parquet::arrow::FileReader * gparquet_arrow_file_reader_get_raw(GParquetArrowFileReader *arrow_file_reader); +parquet::ReaderProperties * +gparquet_reader_properties_get_raw(GParquetReaderProperties *properties); +parquet::ArrowReaderProperties * +gparquet_reader_properties_get_arrow_raw(GParquetReaderProperties *properties); diff --git a/c_glib/test/parquet/test-arrow-file-reader.rb b/c_glib/test/parquet/test-arrow-file-reader.rb index eff5ad966aea..a088f91520bd 100644 --- a/c_glib/test/parquet/test-arrow-file-reader.rb +++ b/c_glib/test/parquet/test-arrow-file-reader.rb @@ -39,6 +39,60 @@ def setup end end + sub_test_case(".new(properties)") do + data("path" => :path, "stream" => :stream) + test("read") do |source_type| + properties = Parquet::ReaderProperties.new + properties.enable_buffered_stream + properties.pre_buffer = false + properties.buffer_size = 4096 + source = if source_type == :path + @file.path + else + Arrow::FileInputStream.new(@file.path) + end + reader = Parquet::ArrowFileReader.new(source, properties) + begin + assert_equal(@table, reader.read_table) + assert_equal(build_table("a" => @a_array.slice(1, 1), + "b" => @b_array.slice(1, 1)), + reader.read_row_group(1)) + ensure + reader.close + reader.unref + source.unref if source_type == :stream + end + end + + test("nonexistent path") do + properties = Parquet::ReaderProperties.new + assert_raise(Arrow::Error::Io) do + Parquet::ArrowFileReader.new("#{@file.path}.nonexistent", properties) + end + end + + data("path" => :path, "stream" => :stream) + test("invalid file") do |source_type| + Tempfile.create("invalid-parquet") do |file| + file.write("not a parquet file") + file.flush + source = if source_type == :path + file.path + else + Arrow::FileInputStream.new(file.path) + end + properties = Parquet::ReaderProperties.new + begin + assert_raise(Arrow::Error::Invalid) do + Parquet::ArrowFileReader.new(source, properties) + end + ensure + source.unref if source_type == :stream + end + end + end + end + def test_schema assert_equal(<<-SCHEMA.chomp, @reader.schema.to_s) a: string diff --git a/c_glib/test/parquet/test-reader-properties.rb b/c_glib/test/parquet/test-reader-properties.rb new file mode 100644 index 000000000000..f1d7598aa30b --- /dev/null +++ b/c_glib/test/parquet/test-reader-properties.rb @@ -0,0 +1,57 @@ +# 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. + +class TestParquetReaderProperties < Test::Unit::TestCase + def setup + omit("Parquet is required") unless defined?(::Parquet) + @properties = Parquet::ReaderProperties.new + end + + def test_buffered_stream + assert do + not @properties.buffered_stream_enabled? + end + @properties.enable_buffered_stream + assert do + @properties.buffered_stream_enabled? + end + @properties.disable_buffered_stream + assert do + not @properties.buffered_stream_enabled? + end + end + + def test_pre_buffer + assert do + @properties.pre_buffer? + end + @properties.pre_buffer = false + assert do + not @properties.pre_buffer? + end + @properties.pre_buffer = true + assert do + @properties.pre_buffer? + end + end + + def test_buffer_size + assert_equal(16 * 1024, @properties.buffer_size) + @properties.buffer_size = 32 * 1024 + assert_equal(32 * 1024, @properties.buffer_size) + end +end