Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 18 additions & 15 deletions documentation/components/extensions/arrow-ext.md
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ The Arrow Rust crates offer additional I/O capabilities that are candidates for
## Features

- Read and write Apache Parquet files through PHP streaming interfaces
- Flat types: INT32, INT64, FLOAT, DOUBLE, BOOLEAN, STRING, BINARY, DATE32, TIMESTAMP
- Flat types: INT32, INT64, FLOAT, DOUBLE, BOOLEAN, STRING, BINARY, DATE, TIMESTAMP, TIME, DECIMAL
- Nested types: LIST, STRUCT, MAP (arbitrarily nested)
- Compression codecs: UNCOMPRESSED, SNAPPY, GZIP, ZSTD, LZ4_RAW, BROTLI
- Column projection for selective reads
Expand Down Expand Up @@ -179,20 +179,23 @@ $writer->close();

The schema is an array of column definitions. Each column has a `name`, `type`, and optional `optional` flag.

| Type | PHP Read Value | Notes |
|------|---------------|-------|
| `BOOLEAN` | `bool` | |
| `INT32` | `int` | |
| `INT64` | `int` | |
| `FLOAT` | `float` | |
| `DOUBLE` | `float` | |
| `STRING` | `string` | |
| `BINARY` | `string` (raw bytes) | |
| `DATE32` | `string` (YYYY-MM-DD) | |
| `TIMESTAMP` | `string` (ISO 8601) | |
| `LIST` | `array` | Requires `children` key with 1 element |
| `STRUCT` | `array` (associative) | Requires `children` key with N elements |
| `MAP` | `array` (associative) | Requires `children` key with 2 elements (key + value) |
| Type | PHP Read Value | Notes |
|-------------|-----------------------|-------------------------------------------------------------------------------------------------------------------|
| `BOOLEAN` | `bool` | |
| `INT32` | `int` | |
| `INT64` | `int` | |
| `FLOAT` | `float` | |
| `DOUBLE` | `float` | |
| `STRING` | `string` | |
| `BINARY` | `string` (raw bytes) | |
| `DATE` | `DateTimeImmutable` | `int` lane: days since epoch |
| `TIMESTAMP` | `DateTimeImmutable` | Keys `unit` (`MILLIS\|MICROS\|NANOS`, default `MICROS`) and `utc` (default `true`); `int` lane in the column unit |
| `TIME` | `DateInterval` | Key `unit` (`MILLIS\|MICROS\|NANOS`, default `MICROS`) |
| `DECIMAL` | `float` | Keys `precision`, `scale` |
| `UUID` | `string` | |
| `LIST` | `array` | Requires `children` key with 1 element |
| `STRUCT` | `array` (associative) | Requires `children` key with N elements |
| `MAP` | `array` (associative) | Requires `children` key with 2 elements (key + value) |

**Nested schema example:**

Expand Down
4 changes: 2 additions & 2 deletions documentation/components/libs/parquet.md
Original file line number Diff line number Diff line change
Expand Up @@ -256,6 +256,8 @@ $schema = Schema::with(
);
```

`FlatColumn::dateTime()` is written as `TIMESTAMP(isAdjustedToUTC=true, MICROS)`; files with `isAdjustedToUTC=false` read as UTC wall clock.

Once we have a schema, we can create a writer.

```php
Expand Down Expand Up @@ -333,8 +335,6 @@ $writer->close();
- `GZIP_COMPRESSION_LEVEL` - default: `9` - compression level for GZIP compression (applied only when GZIP compression
is enabled).
- `PAGE_SIZE_BYTES` - default: `8Kb` - maximum size of data page.
- `ROUND_NANOSECONDS` - default: `false` - Since PHP does not support nanoseconds precision for DateTime objects, when
this options is set to true, reader will round nanoseconds to microseconds.
- `ROW_GROUP_SIZE_BYTES` - default: `8Mb` - maximum size of row group.
- `ROW_GROUP_SIZE_CHECK_INTERVAL` default: `1000` - number of rows to write before checking if row group size limit is
reached.
Expand Down
48 changes: 48 additions & 0 deletions documentation/upgrading.md
Original file line number Diff line number Diff line change
Expand Up @@ -419,6 +419,54 @@ The same applies to map values and structure elements.
| `new DateTimeDefinition('at', true)` | `new DateTimeDefinition('at', type_datetime(), true)` or `datetime_schema('at', nullable: true, zone: 'Europe/Warsaw')` |
| a `flow_php` extension older than this release | ignored - the PHP engine runs; a current extension with an older `flow-php/etl` keeps that library's behaviour |

### 41) `flow-php/parquet` - `Converter::isFor()` replaced by `static Converter::forColumn()`, `Int32DateTimeConverter` removed

| Before | After |
|--------------------------------------------------------------------|------------------------------------------------------------------------------------|
| `Converter::isFor(FlatColumn, Options): bool` on a shared instance | `static Converter::forColumn(FlatColumn, Options): ?self` - a converter per column |
| `new DataConverter([new TimeConverter(), ...], $options)` | `new DataConverter([TimeConverter::class, ...], $options)` |
| `Int32DateTimeConverter` | removed - INT32 TIMESTAMP is not a legal Parquet carrier |

### 42) `flow-php/parquet` - `LogicalType\Timestamp` / `Time` take a `TimeUnit`

| Before | After |
|-------------------------------------------------------------|--------------------------------------------------------------------------------------------------|
| `new Timestamp($isAdjustedToUTC, $millis, $micros, $nanos)` | `new Timestamp($isAdjustedToUTC, TimeUnit::MICROSECONDS)`, same for `Time` |
| `TimeUnit` - pure enum with `MICROSECONDS` only | string-backed enum `MILLISECONDS = 'MILLIS'`, `MICROSECONDS = 'MICROS'`, `NANOSECONDS = 'NANOS'` |

### 43) `flow-php/parquet` - `Option::ROUND_NANOSECONDS` removed

| Before | After |
|-----------------------------|---------------------------------------------------------------------------------------|
| `Option::ROUND_NANOSECONDS` | removed - NANOS timestamps always read as `DateTimeImmutable` floored to microseconds |

### 44) `flow-php/parquet` - `encode_decimal()` / `decode_decimal()` drop `ByteOrder` and the read-side precision check

| Before | After |
|---------------------------------------------------------------------------------|--------------------------------------------------------------------------------------------------------------|
| `encode_decimal(ByteOrder, float, int $byteLength, int $precision, int $scale)` | `encode_decimal(float, int $precision, int $scale, ?int $byteLength)` - `null` = minimal length (BYTE_ARRAY) |
| `decode_decimal(ByteOrder, string, int $precision, int $scale)` | `decode_decimal(string, int $scale)` - no precision check |

### 45) `flow-php/parquet`, `flow-php/arrow-ext` - TIMESTAMP columns are written `isAdjustedToUTC=true`

| Before | After |
|--------------------------------------------------------------------------------------------------------------|------------------------------------------------------------------------------------------------------|
| `FlatColumn::dateTime()` / arrow `TIMESTAMP` written `isAdjustedToUTC=false` (DuckDB/Spark read `TIMESTAMP`) | written `isAdjustedToUTC=true` (DuckDB/Spark read `TIMESTAMPTZ`); reading `false` files is unchanged |

### 46) `flow-php/parquet`, `flow-php/etl-adapter-parquet` - ConvertedType-only TIMESTAMP/TIME/DECIMAL columns are typed

| Before | After |
|-------------------------------------------------------------------------------|------------|
| `TIMESTAMP_MILLIS` / `TIMESTAMP_MICROS` column without a logical type → `int` | `datetime` |
| `TIME_MILLIS` / `TIME_MICROS` column without a logical type → `int` | `time` |
| `DECIMAL` converted type on INT32/INT64 → unscaled `int` | `float` |

### 47) `flow-php/arrow-ext` - `Writer` DATE `int` lane is days since epoch

| Before | After |
|---------------------------------------------|-------------------------------------|
| `'d' => [1641600000]` - seconds since epoch | `'d' => [19000]` - days since epoch |

---

## Upgrading from 0.43.x to 0.44.x
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
<?php

declare(strict_types=1);

namespace Flow\ETL\Adapter\Parquet\Tests\Integration;

use DateInterval;
use DateTimeInterface;
use Flow\ETL\Tests\FlowTestCase;
use Flow\Parquet\Engine\PhpParquetEngine;
use Flow\Parquet\ParquetFile\Schema;
use Flow\Parquet\ParquetFile\Schema\ConvertedType;
use Flow\Parquet\ParquetFile\Schema\FlatColumn;
use Flow\Parquet\ParquetFile\Schema\PhysicalType;
use Flow\Parquet\ParquetFile\Schema\Repetition;
use Flow\Parquet\Writer;

use function Flow\ETL\Adapter\Parquet\from_parquet;
use function Flow\ETL\DSL\data_frame;
use function Flow\ETL\DSL\datetime_schema;
use function Flow\ETL\DSL\float_schema;
use function Flow\ETL\DSL\schema;
use function Flow\ETL\DSL\time_schema;
use function Flow\Filesystem\DSL\memory_filesystem;
use function Flow\Filesystem\DSL\path;

final class ConvertedTypeOnlyColumnsTest extends FlowTestCase
{
public function test_converted_type_only_columns_read_as_temporal_and_float_entries(): void
{
$memory = memory_filesystem();
$path = path('memory://var/converted_type_only_columns.parquet');
$writer = new Writer(engine: new PhpParquetEngine());

$writer->openForStream(
$memory->writeTo($path),
Schema::with(
new FlatColumn(
'ts_ms',
PhysicalType::INT64,
ConvertedType::TIMESTAMP_MILLIS,
null,
Repetition::OPTIONAL,
),
new FlatColumn(
'ts_us',
PhysicalType::INT64,
ConvertedType::TIMESTAMP_MICROS,
null,
Repetition::OPTIONAL,
),
new FlatColumn('t_ms', PhysicalType::INT32, ConvertedType::TIME_MILLIS, null, Repetition::OPTIONAL),
new FlatColumn('t_us', PhysicalType::INT64, ConvertedType::TIME_MICROS, null, Repetition::OPTIONAL),
new FlatColumn('dec9', PhysicalType::INT32, ConvertedType::DECIMAL, null, Repetition::OPTIONAL, 9, 2),
),
);
$writer->writeBatch([[
'ts_ms' => 1577934245678,
'ts_us' => 1577934245678901,
't_ms' => 11045678,
't_us' => 11045678901,
'dec9' => 1234567,
]]);
$writer->close();

$rows = data_frame()->read(from_parquet($path, filesystem: $memory))->fetch();

static::assertEquals(
schema(
datetime_schema('ts_ms', true),
datetime_schema('ts_us', true),
time_schema('t_ms', true),
time_schema('t_us', true),
float_schema('dec9', true),
),
$rows->schema(),
);

$row = $rows[0];
$tsMs = $row->get('ts_ms');
$tsUs = $row->get('ts_us');
$tMs = $row->get('t_ms');
$tUs = $row->get('t_us');

static::assertInstanceOf(DateTimeInterface::class, $tsMs);
static::assertInstanceOf(DateTimeInterface::class, $tsUs);
static::assertInstanceOf(DateInterval::class, $tMs);
static::assertInstanceOf(DateInterval::class, $tUs);
static::assertSame('2020-01-02 03:04:05.678000 +00:00', $tsMs->format('Y-m-d H:i:s.u P'));
static::assertSame('2020-01-02 03:04:05.678901 +00:00', $tsUs->format('Y-m-d H:i:s.u P'));
static::assertSame('03:04:05.678000', $tMs->format('%H:%I:%S.%F'));
static::assertSame('03:04:05.678901', $tUs->format('%H:%I:%S.%F'));
static::assertSame(12345.67, $row->get('dec9'));
}
}
28 changes: 25 additions & 3 deletions src/extension/arrow-ext/src/parquet/reader.rs
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
use arrow_schema::DataType;
use arrow_schema::{DataType, TimeUnit};
use ext_php_rs::boxed::ZBox;
use ext_php_rs::prelude::*;
use ext_php_rs::types::{ZendHashTable, ZendObject};
Expand Down Expand Up @@ -108,6 +108,15 @@ impl Reader {
}
}

fn time_unit_name(unit: &TimeUnit) -> &'static str {
match unit {
// Parquet has no second unit; arrow-rs stores seconds as MILLIS.
TimeUnit::Second | TimeUnit::Millisecond => "MILLIS",
TimeUnit::Microsecond => "MICROS",
TimeUnit::Nanosecond => "NANOS",
}
}

fn field_to_php_schema(field: &arrow_schema::Field) -> PhpResult<ZBox<ZendHashTable>> {
let mut entry = ZendHashTable::new();
entry
Expand All @@ -131,7 +140,7 @@ fn field_to_php_schema(field: &arrow_schema::Field) -> PhpResult<ZBox<ZendHashTa
DataType::Date32 | DataType::Date64 => "DATE",
DataType::Timestamp(_, _) => "TIMESTAMP",
DataType::Time32(_) | DataType::Time64(_) => "TIME",
DataType::Decimal128(_, _) => "DECIMAL",
DataType::Decimal128(_, _) | DataType::Decimal256(_, _) => "DECIMAL",
DataType::FixedSizeBinary(_) => "FIXED_SIZE_BINARY",
DataType::List(_) | DataType::LargeList(_) => "LIST",
DataType::Struct(_) => "STRUCT",
Expand All @@ -152,14 +161,27 @@ fn field_to_php_schema(field: &arrow_schema::Field) -> PhpResult<ZBox<ZendHashTa
.map_err(|_| parquet_exception("Failed to build schema entry"))?;

match field.data_type() {
DataType::Decimal128(precision, scale) => {
DataType::Decimal128(precision, scale) | DataType::Decimal256(precision, scale) => {
entry
.insert("precision", *precision as i64)
.map_err(|_| parquet_exception("Failed to build schema entry"))?;
entry
.insert("scale", *scale as i64)
.map_err(|_| parquet_exception("Failed to build schema entry"))?;
}
DataType::Timestamp(unit, tz) => {
entry
.insert("unit", time_unit_name(unit))
.map_err(|_| parquet_exception("Failed to build schema entry"))?;
entry
.insert("utc", tz.is_some())
.map_err(|_| parquet_exception("Failed to build schema entry"))?;
}
DataType::Time32(unit) | DataType::Time64(unit) => {
entry
.insert("unit", time_unit_name(unit))
.map_err(|_| parquet_exception("Failed to build schema entry"))?;
}
DataType::FixedSizeBinary(n) => {
entry
.insert("length", *n as i64)
Expand Down
Loading
Loading