blob: d190cc29923e8ee5c4d54542ed14b5bb0385e596 [file]
# 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.
require "bigdecimal"
require_relative "array-builder"
require_relative "bitmap"
require_relative "bitmap-builder"
module ArrowFormat
using FlatBuffers::AppendAsBytes if FlatBuffers.const_defined?(:AppendAsBytes)
class Array
include Enumerable
class << self
def build(values)
ArrayBuilder.new(values).build
end
end
attr_reader :type
attr_reader :size
alias_method :length, :size
attr_reader :offset
attr_reader :validity_buffer
def initialize(type, size, validity_buffer)
@type = type
@size = size
@offset = 0
@validity_buffer = validity_buffer
@sliced_buffers = {}
end
def slice(offset, size=nil)
sliced = dup
sliced.slice!(@offset + offset, size || @size - offset)
sliced
end
def valid?(i)
return true if @validity_buffer.nil?
validity_bitmap[i]
end
def null?(i)
not valid?(i)
end
def [](i)
if valid?(i)
unpack_value(i)
else
nil
end
end
def n_nulls
if @validity_buffer.nil?
0
else
@size - validity_bitmap.popcount
end
end
def empty?
@size.zero?
end
def ==(other)
return false unless other.is_a?(self.class)
return false unless @size == other.size
return true if @validity_buffer.nil? and other.validity_buffer.nil?
if @offset == other.offset and @validity_buffer == other.validity_buffer
return true
end
validity_bitmap == other.validity_bitmap
end
protected
def slice!(offset, size)
@offset = offset
@size = size
clear_cache
end
def validity_bitmap
@validity_bitmap ||= Bitmap.new(@validity_buffer, @offset, @size)
end
private
def apply_validity(array)
return array if @validity_buffer.nil?
validity_bitmap.each_with_index do |is_valid, i|
array[i] = nil unless is_valid
end
array
end
def clear_cache
@validity_bitmap = nil
@sliced_buffers = {}
end
def slice_buffer(id, buffer)
return buffer if buffer.nil?
return buffer if @offset.zero?
@sliced_buffers[id] ||= yield(buffer)
end
def slice_bitmap_buffer(id, buffer)
slice_buffer(id, buffer) do
if (@offset % 8).zero?
buffer.slice(@offset / 8)
else
# We need to copy because we can't do bit level slice.
# TODO: Optimize.
valid_bytes = []
Bitmap.new(buffer, @offset, @size).each_slice(8) do |valids|
valid_byte = 0
valids.each_with_index do |valid, i|
valid_byte |= 1 << (i % 8) if valid
end
valid_bytes << valid_byte
end
IO::Buffer.for(valid_bytes.pack("C*"))
end
end
end
def slice_fixed_element_size_buffer(id, buffer, element_size)
slice_buffer(id, buffer) do
buffer.slice(element_size * @offset)
end
end
def slice_offsets_buffer(id, buffer, buffer_type)
slice_buffer(id, buffer) do
offset_size = IO::Buffer.size_of(buffer_type)
buffer_offset = offset_size * @offset
first_offset = nil
# TODO: Optimize
sliced_buffer = IO::Buffer.new(offset_size * (@size + 1))
buffer.each(buffer_type,
buffer_offset,
@size + 1).with_index do |(_, offset), i|
first_offset ||= offset
new_offset = offset - first_offset
sliced_buffer.set_value(buffer_type,
offset_size * i,
new_offset)
end
sliced_buffer
end
end
end
class NullArray < Array
def initialize(size)
super(NullType.singleton, size, nil)
end
def valid?(i)
false
end
def null?(i)
true
end
def each_buffer
return to_enum(__method__) unless block_given?
end
def n_nulls
@size
end
def to_a
[nil] * @size
end
def each
return to_enum(__method__) unless block_given?
@size.times do
yield(nil)
end
end
end
class PrimitiveArray < Array
include BufferAlignable
attr_reader :values_buffer
def initialize(*args)
n_args = args.size
if self.class.respond_to?(:type)
type = self.class.type
expected_n_args = "1 or 3"
else
type = args.shift
unless type.is_a?(Type)
type = self.class.type_class.try_convert(type) || type
end
expected_n_args = "2 or 4"
end
args = build_data(args[0], type) if args.size == 1
if args.size != 3
message =
"wrong number of arguments " +
"(given #{n_args}, expected #{expected_n_args})"
raise ArgumentError, message
end
size, validity_buffer, values_buffer = args
super(type, size, validity_buffer)
@values_buffer = values_buffer
end
def to_a
return [] if empty?
offset = element_size * @offset
apply_validity(@values_buffer.values(@type.buffer_type, offset, @size))
end
def each(&block)
return to_enum(__method__) {@size} unless block_given?
each_value = Enumerator.new(@size) do |yielder|
offset = element_size * @offset
@values_buffer.each(@type.buffer_type, offset, @size) do |_, value|
yielder << value
end
end
if @validity_buffer.nil?
each_value.each(&block)
else
validity_bitmap.zip(each_value) do |is_valid, value|
if is_valid
yield(value)
else
yield(nil)
end
end
end
end
def each_buffer
return to_enum(__method__) {2} unless block_given?
yield(slice_bitmap_buffer(:validity, @validity_buffer))
yield(slice_fixed_element_size_buffer(:values,
@values_buffer,
element_size))
end
def ==(other)
return false unless super(other)
if @offset == other.offset and @values_buffer == other.values_buffer
return true
end
lazy.zip(other).all? do |value, other_value|
value == other_value
end
end
private
def element_size
IO::Buffer.size_of(@type.buffer_type)
end
def build_data(data, type)
n = 0
validity_buffer_builder = nil
buffer = +"".b
pack_template = type.pack_template
data.each_with_index do |value, i|
if value.nil?
validity_buffer_builder ||= SparseBitmapBuilder.new
validity_buffer_builder.unset(i)
buffer.append_as_bytes(pack_value(nil, pack_template, type))
else
buffer.append_as_bytes(pack_value(value, pack_template, type))
end
n += 1
end
validity_buffer = validity_buffer_builder&.finish(n)
pad!(buffer, buffer_padding_size(buffer))
buffer.freeze
return n, validity_buffer, IO::Buffer.for(buffer)
end
def pack_value(value, template, type)
value = 0 if value.nil?
[value].pack(template)
end
def unpack_value(i)
offset = element_size * (@offset + i)
@values_buffer.get_value(@type.buffer_type, offset)
end
end
class BooleanArray < PrimitiveArray
class << self
def type
BooleanType.singleton
end
end
def to_a
return [] if empty?
values = values_bitmap.to_a
apply_validity(values)
end
def each(&block)
return to_enum(__method__) {@size} unless block_given?
if @validity_buffer.nil?
values_bitmap.each(&block)
else
validity_bitmap.zip(values_bitmap) do |is_valid, value|
if is_valid
yield(value)
else
yield(nil)
end
end
end
end
def each_buffer
return to_enum(__method__) {2} unless block_given?
yield(slice_bitmap_buffer(:validity, @validity_buffer))
yield(slice_bitmap_buffer(:values, @values_buffer))
end
private
def values_bitmap
@values_bitmap ||= Bitmap.new(@values_buffer, @offset, @size)
end
def clear_cache
super
@values_bitmap = nil
end
def build_data(data, type)
n = 0
validity_buffer_builder = nil
values_buffer_builder = DenseBitmapBuilder.new
data.each_with_index do |value, i|
if value.nil?
validity_buffer_builder ||= SparseBitmapBuilder.new
validity_buffer_builder.unset(i)
values_buffer_builder.append(false)
elsif value
values_buffer_builder.append(true)
else
values_buffer_builder.append(false)
end
n += 1
end
validity_buffer = validity_buffer_builder&.finish(n)
return n, validity_buffer, values_buffer_builder.finish
end
def unpack_value(i)
values_bitmap[i]
end
end
class IntArray < PrimitiveArray
end
class Int8Array < IntArray
class << self
def type
Int8Type.singleton
end
end
end
class UInt8Array < IntArray
class << self
def type
UInt8Type.singleton
end
end
end
class Int16Array < IntArray
class << self
def type
Int16Type.singleton
end
end
end
class UInt16Array < IntArray
class << self
def type
UInt16Type.singleton
end
end
end
class Int32Array < IntArray
class << self
def type
Int32Type.singleton
end
end
end
class UInt32Array < IntArray
class << self
def type
UInt32Type.singleton
end
end
end
class Int64Array < IntArray
class << self
def type
Int64Type.singleton
end
end
end
class UInt64Array < IntArray
class << self
def type
UInt64Type.singleton
end
end
end
class FloatingPointArray < PrimitiveArray
end
class Float32Array < FloatingPointArray
class << self
def type
Float32Type.singleton
end
end
end
class Float64Array < FloatingPointArray
class << self
def type
Float64Type.singleton
end
end
end
class TemporalArray < PrimitiveArray
end
class DateArray < TemporalArray
end
class Date32Array < DateArray
class << self
def type
Date32Type.singleton
end
end
private
def pack_value(value, template, type)
if value.nil?
[0].pack(template)
elsif value.is_a?(Date)
[value.day].pack(template)
else
[value].pack(template)
end
end
end
class Date64Array < DateArray
class << self
def type
Date64Type.singleton
end
end
end
class TimeArray < TemporalArray
end
class Time32Array < TimeArray
class << self
def type_class
Time32Type
end
end
end
class Time64Array < TimeArray
class << self
def type_class
Time64Type
end
end
end
class TimestampArray < TemporalArray
class << self
def type_class
TimestampType
end
end
private
def pack_value(value, template, type)
if value.nil?
[0, 0].pack(template)
elsif value.is_a?(Time)
case type.unit
when :second
value = value.to_i
when :millisecond
value = (value.to_i * 1_000) + (value.nsec / 1_000_000)
when :microsecond
value = (value.to_i * 1_000_000) + (value.nsec / 1_000)
when :nanosecond
value = (value.to_i * 1_000_000_000) + value.nsec
end
[value].pack(template)
else
[value].pack(template)
end
end
end
class IntervalArray < TemporalArray
end
class YearMonthIntervalArray < IntervalArray
class << self
def type
YearMonthIntervalType.singleton
end
end
end
class DayTimeIntervalArray < IntervalArray
class << self
def type
DayTimeIntervalType.singleton
end
end
def to_a
return [] if empty?
offset = element_size * @offset
values = @values_buffer.
each(@type.buffer_type, offset, @size * 2).
each_slice(2).
collect do |(_, day), (_, millisecond)|
[day, millisecond]
end
apply_validity(values)
end
def each(&block)
return to_enum(__method__) {@size} unless block_given?
each_value = Enumerator.new(@size) do |yielder|
offset = element_size * @offset
@values_buffer.
each(@type.buffer_type, offset, @size * 2).
each_slice(2) do |(_, day), (_, millisecond)|
yielder << [day, millisecond]
end
end
if @validity_buffer.nil?
each_value.each(&block)
else
validity_bitmap.zip(each_value) do |is_valid, value|
if is_valid
yield(value)
else
yield(nil)
end
end
end
end
private
def element_size
super * 2
end
def pack_value(value, template, type)
if value.nil?
[0, 0].pack(template)
elsif value.is_a?(Hash)
[value[:day], value[:millisecond]].pack(template)
else
value.pack(template)
end
end
def unpack_value(i)
offset = element_size * (@offset + i)
@values_buffer.get_values([@type.buffer_type] * 2, offset)
end
end
class MonthDayNanoIntervalArray < IntervalArray
class << self
def type
MonthDayNanoIntervalType.singleton
end
end
def to_a
return [] if empty?
buffer_types = @type.buffer_types
value_size = IO::Buffer.size_of(buffer_types)
base_offset = value_size * @offset
values = @size.times.collect do |i|
offset = base_offset + value_size * i
@values_buffer.get_values(buffer_types, offset)
end
apply_validity(values)
end
def each(&block)
return to_enum(__method__) {@size} unless block_given?
each_value = Enumerator.new(@size) do |yielder|
buffer_types = @type.buffer_types
value_size = IO::Buffer.size_of(buffer_types)
base_offset = value_size * @offset
@size.times do |i|
offset = base_offset + value_size * i
yielder << @values_buffer.get_values(buffer_types, offset)
end
end
if @validity_buffer.nil?
each_value.each(&block)
else
validity_bitmap.zip(each_value) do |is_valid, value|
if is_valid
yield(value)
else
yield(nil)
end
end
end
end
private
def element_size
IO::Buffer.size_of(@type.buffer_types)
end
def pack_value(value, template, type)
if value.nil?
[0, 0, 0].pack(template)
elsif value.is_a?(Hash)
[value[:month], value[:day], value[:nanosecond]].pack(template)
else
value.pack(template)
end
end
def unpack_value(i)
buffer_types = @type.buffer_types
value_size = IO::Buffer.size_of(buffer_types)
offset = value_size * (@offset + i)
@values_buffer.get_values(buffer_types, offset)
end
end
class DurationArray < TemporalArray
class << self
def type_class
DurationType
end
end
end
class VariableSizeBinaryArray < Array
include BufferAlignable
def initialize(*args)
if args.size == 1
args = build_data(args.first)
elsif args.size != 4
raise ArgumentError,
"wrong number of arguments (given #{args.size}, expected 1 or 4)"
end
size, validity_buffer, offsets_buffer, values_buffer = args
super(self.class.type, size, validity_buffer)
@offsets_buffer = offsets_buffer
@values_buffer = values_buffer
end
def each_buffer
return to_enum(__method__) unless block_given?
yield(slice_bitmap_buffer(:validity, @validity_buffer))
yield(slice_offsets_buffer(:offsets,
@offsets_buffer,
@type.offset_buffer_type))
sliced_values_buffer = slice_buffer(:values, @values_buffer) do
first_offset = @offsets_buffer.get_value(@type.offset_buffer_type,
offset_size * @offset)
@values_buffer.slice(first_offset)
end
yield(sliced_values_buffer)
end
def offsets
return [0] if empty?
@offsets_buffer.
each(@type.offset_buffer_type, offset_size * @offset, @size + 1).
collect {|_, offset| offset}
end
def to_a
return [] if empty?
values = offsets.each_cons(2).collect do |offset, next_offset|
length = next_offset - offset
@values_buffer.get_string(offset, length, @type.encoding)
end
apply_validity(values)
end
private
def build_data(data)
type = self.class.type
n = 0
validity_buffer_builder = nil
values = +"".b
offsets = [0]
data.each_with_index do |value, i|
if value.nil?
validity_buffer_builder ||= SparseBitmapBuilder.new
validity_buffer_builder.unset(i)
offsets << values.bytesize
else
values.append_as_bytes(value)
offsets << values.bytesize
end
n += 1
end
validity_buffer = validity_buffer_builder&.finish(n)
offsets_data = offsets.pack("#{type.offset_pack_template}*")
pad!(offsets_data, buffer_padding_size(offsets_data))
offsets_data.freeze
offsets_buffer = IO::Buffer.for(offsets_data)
pad!(values, buffer_padding_size(values))
values.freeze
values_buffer = IO::Buffer.for(values)
[
n,
validity_buffer,
offsets_buffer,
values_buffer,
]
end
def offset_size
IO::Buffer.size_of(@type.offset_buffer_type)
end
def unpack_value(i)
offset_types = [@type.offset_buffer_type] * 2
offsets_offset = offset_size * (@offset + i)
offset, next_offset = @offsets_buffer.get_values(offset_types,
offsets_offset)
length = next_offset - offset
@values_buffer.get_string(offset, length, @type.encoding)
end
end
class BinaryArray < VariableSizeBinaryArray
class << self
def type
BinaryType.singleton
end
end
end
class LargeBinaryArray < VariableSizeBinaryArray
class << self
def type
LargeBinaryType.singleton
end
end
end
class VariableSizeUTF8Array < VariableSizeBinaryArray
end
class UTF8Array < VariableSizeUTF8Array
class << self
def type
UTF8Type.singleton
end
end
end
class LargeUTF8Array < VariableSizeUTF8Array
class << self
def type
LargeUTF8Type.singleton
end
end
end
class FixedSizeBinaryArray < Array
def initialize(type, size, validity_buffer, values_buffer)
super(type, size, validity_buffer)
@values_buffer = values_buffer
end
def each_buffer
return to_enum(__method__) unless block_given?
yield(slice_bitmap_buffer(:validity, @validity_buffer))
yield(slice_fixed_element_size_buffer(:values,
@values_buffer,
@type.byte_width))
end
def to_a
return [] if empty?
byte_width = @type.byte_width
values = 0.step(@size * byte_width - 1, byte_width).collect do |offset|
@values_buffer.get_string(offset, byte_width)
end
apply_validity(values)
end
end
class DecimalArray < FixedSizeBinaryArray
def to_a
return [] if empty?
byte_width = @type.byte_width
buffer_types = [:u64] * (byte_width / 8 - 1) + [:s64]
base_offset = byte_width * @offset
values = 0.step(@size * byte_width - 1, byte_width).collect do |offset|
@values_buffer.get_values(buffer_types, base_offset + offset)
end
apply_validity(values).collect do |value|
if value.nil?
nil
else
BigDecimal(format_value(value))
end
end
end
private
def format_value(components)
highest = components.last
width = @type.precision
width += 1 if highest < 0
value = 0
components.reverse_each do |component|
value = (value << 64) + component
end
string = value.to_s
if @type.scale < 0
string << ("0" * -@type.scale)
elsif @type.scale > 0
n_digits = string.bytesize
n_digits -= 1 if value < 0
if n_digits <= @type.scale
prefix = "0." + ("0" * (@type.scale - n_digits))
if value < 0
string[1, 0] = prefix
else
string[0, 0] = prefix
end
else
string[-@type.scale, 0] = "."
end
end
string
end
end
class Decimal128Array < DecimalArray
end
class Decimal256Array < DecimalArray
end
class VariableSizeListArray < Array
attr_reader :child
def initialize(type, size, validity_buffer, offsets_buffer, child)
super(type, size, validity_buffer)
@offsets_buffer = offsets_buffer
@child = child
end
def each_buffer(&block)
return to_enum(__method__) unless block_given?
yield(slice_bitmap_buffer(:validity, @validity_buffer))
yield(slice_offsets_buffer(:offsets,
@offsets_buffer,
@type.offset_buffer_type))
end
def offsets
return [0] if empty?
@offsets_buffer.
each(@type.offset_buffer_type, offset_size * @offset, @size + 1).
collect {|_, offset| offset}
end
def to_a
return [] if empty?
child_values = @child.to_a
values = offsets.each_cons(2).collect do |offset, next_offset|
child_values[offset...next_offset]
end
apply_validity(values)
end
private
def offset_size
IO::Buffer.size_of(@type.offset_buffer_type)
end
def slice!(offset, size)
super
first_offset =
@offsets_buffer.get_value(@type.offset_buffer_type,
offset_size * @offset)
last_offset =
@offsets_buffer.get_value(@type.offset_buffer_type,
offset_size * (@offset + @size + 1))
@child = @child.slice(first_offset, last_offset - first_offset)
end
end
class ListArray < VariableSizeListArray
end
class LargeListArray < VariableSizeListArray
end
class FixedSizeListArray < Array
attr_reader :child
def initialize(type, size, validity_buffer, child)
super(type, size, validity_buffer)
@child = child
end
def each_buffer(&block)
return to_enum(__method__) unless block_given?
yield(slice_bitmap_buffer(:validity, @validity_buffer))
end
def to_a
return [] if empty?
values = @child.to_a.each_slice(@type.size).to_a
apply_validity(values)
end
private
def slice!(offset, size)
super
@child = @child.slice(@type.size * @offset,
@type.size * (@offset + @size + 1))
end
end
class StructArray < Array
attr_reader :children
def initialize(type, size, validity_buffer, children)
super(type, size, validity_buffer)
@children = children
end
def each_buffer(&block)
return to_enum(__method__) unless block_given?
yield(slice_bitmap_buffer(:validity, @validity_buffer))
end
def to_a
return [] if empty?
if @children.empty?
values = [[]] * @size
else
children_values = @children.collect(&:to_a)
values = children_values[0].zip(*children_values[1..-1])
end
apply_validity(values)
end
private
def slice!(offset, size)
super
@children = @children.collect do |child|
child.slice(offset, size)
end
end
end
class MapArray < VariableSizeListArray
def to_a
return [] if empty?
super.collect do |entries|
if entries.nil?
entries
else
hash = {}
entries.each do |key, value|
hash[key] = value
end
hash
end
end
end
end
class UnionArray < Array
attr_reader :types_buffer
attr_reader :children
def initialize(type, size, types_buffer, children)
super(type, size, nil)
@types_buffer = types_buffer
@children = children
end
def each_type(&block)
return [].each(&block) if empty?
return to_enum(__method__) unless block_given?
@types_buffer.each(type_buffer_type,
type_element_size * @offset,
@size) do |_, type|
yield(type)
end
end
private
def type_buffer_type
:S8
end
def type_element_size
IO::Buffer.size_of(type_buffer_type)
end
end
class DenseUnionArray < UnionArray
def initialize(type,
size,
types_buffer,
offsets_buffer,
children)
super(type, size, types_buffer, children)
@offsets_buffer = offsets_buffer
end
def each_buffer(&block)
return to_enum(__method__) unless block_given?
# TODO: Dictionary delta support (slice support)
yield(@types_buffer)
yield(@offsets_buffer)
end
def each_offset(&block)
return [].each(&block) if empty?
return to_enum(__method__) unless block_given?
@offsets_buffer.each(@type.offset_buffer_type,
offset_element_size * @offset,
@size) do |_, offset|
yield(offset)
end
end
def to_a
return [] if empty?
children_values = @children.collect(&:to_a)
each_type.zip(each_offset).collect do |type, offset|
index = @type.resolve_type_index(type)
children_values[index][offset]
end
end
private
def offset_buffer_type
:s32
end
def offset_element_size
IO::Buffer.size_of(offset_buffer_type)
end
end
class SparseUnionArray < UnionArray
def each_buffer(&block)
return to_enum(__method__) unless block_given?
yield(slice_fixed_element_size_buffer(:types,
@types_buffer,
type_element_size))
end
def to_a
return [] if empty?
children_values = @children.collect(&:to_a)
each_type.with_index.collect do |type, i|
index = @type.resolve_type_index(type)
children_values[index][i]
end
end
private
def slice!(offset, size)
super
@children = @children.collect do |child|
child.slice(offset, size)
end
end
end
class DictionaryArray < Array
attr_reader :indices_buffer
attr_reader :dictionaries
def initialize(type,
size,
validity_buffer,
indices_buffer,
dictionaries)
super(type, size, validity_buffer)
@indices_buffer = indices_buffer
@dictionaries = dictionaries
end
# TODO: Slice support
def each_buffer
return to_enum(__method__) unless block_given?
yield(@validity_buffer)
yield(@indices_buffer)
end
def indices
buffer_type = @type.index_type.buffer_type
offset = IO::Buffer.size_of(buffer_type) * @offset
apply_validity(@indices_buffer.values(buffer_type, offset, @size))
end
def to_a
return [] if empty?
values = []
@dictionaries.each do |dictionary|
values.concat(dictionary.array.to_a)
end
indices.collect do |index|
if index.nil?
nil
else
values[index]
end
end
end
end
end