Writes data to Parquet files with support for chunking, consolidation, Hive partitioning, and automatic object store uploads.
ParquetOutput writes data to Parquet files with support for chunking, consolidation, Hive partitioning, and automatic object store uploads. It inherits from the base Output class and provides specialized functionality for columnar data storage.
ParquetOutput
Classapplication_sdk.outputs.parquetOutputWrites data to Parquet files with support for chunking, consolidation, Hive partitioning, and automatic object store uploads.
Methods5
__init__
__init__(self, output_path: str, output_suffix: str = '', typename: Optional[str] = None, chunk_size: Optional[int] = 100000, buffer_size: int = 5000, chunk_start: Optional[int] = None, start_marker: Optional[str] = None, end_marker: Optional[str] = None, retain_local_copy: bool = False, use_consolidation: bool = False)Parameters
output_pathstroutput_suffixstrtypenameOptional[str]chunk_sizeOptional[int]buffer_sizeintchunk_startOptional[int]start_markerOptional[str]end_markerOptional[str]retain_local_copybooluse_consolidationboolwrite_dataframe
async write_dataframe(self, dataframe: pd.DataFrame) -> NoneParameters
dataframepd.DataFramewrite_batched_dataframe
async write_batched_dataframe(self, batch_df: pd.DataFrame) -> NoneParameters
batch_dfpd.DataFramewrite_daft_dataframe
async write_daft_dataframe(self, dataframe: daft.DataFrame, partition_cols: Optional[List] = None, write_mode: Union[WriteMode, str] = WriteMode.APPEND, morsel_size: int = 100000) -> NoneParameters
dataframedaft.DataFramepartition_colsOptional[List]write_modeUnion[WriteMode, str]morsel_sizeintget_full_path
get_full_path(self) -> strReturns
str - Full output directory pathUsage Examples
Basic initialization
Initialize ParquetOutput with basic configuration
from application_sdk.outputs import ParquetOutput
parquet_output = ParquetOutput(
output_path="/tmp/output",
output_suffix="data",
typename="table",
chunk_size=100000,
buffer_size=5000
)
await parquet_output.write_dataframe(df)
With consolidation
Use consolidation for efficient buffered writing of large datasets
parquet_output = ParquetOutput(
output_path="/tmp/output",
output_suffix="data",
use_consolidation=True,
chunk_size=100000
)
await parquet_output.write_batched_dataframe(batched_df)
With Hive partitioning
Write with Hive-style partitioning for efficient data organization
from application_sdk.outputs import ParquetOutput, WriteMode
import daft
parquet_output = ParquetOutput(
output_path="/tmp/output",
output_suffix="partitioned_data"
)
await parquet_output.write_daft_dataframe(
daft_df,
partition_cols=["year", "month"],
write_mode=WriteMode.OVERWRITE
)
Usage patterns
For detailed usage patterns including Hive partitioning, consolidation, and other ParquetOutput-specific features, see Output usage patterns and select the ParquetOutput tab.
See also
- Outputs: Base Output class and common usage patterns for all output types
- JsonOutput: Write data to JSON files in JSONL format
- IcebergOutput: Write data to Apache Iceberg tables with automatic table creation