|
148 | 148 | ) |
149 | 149 | from pyiceberg.table import DOWNCAST_NS_TIMESTAMP_TO_US_ON_WRITE, TableProperties |
150 | 150 | from pyiceberg.table.deletion_vector import deletion_vectors_from_puffin_file |
151 | | -from pyiceberg.table.locations import load_location_provider |
| 151 | +from pyiceberg.table.locations import LocationProvider, load_location_provider |
152 | 152 | from pyiceberg.table.metadata import TableMetadata |
153 | 153 | from pyiceberg.table.name_mapping import NameMapping, apply_name_mapping |
154 | 154 | from pyiceberg.table.puffin import PuffinFile |
@@ -2791,25 +2791,24 @@ def _build_data_file( |
2791 | 2791 | ) |
2792 | 2792 |
|
2793 | 2793 |
|
2794 | | -def write_file(io: FileIO, table_metadata: TableMetadata, tasks: Iterator[WriteTask]) -> Iterator[DataFile]: |
2795 | | - from pyiceberg.table import DOWNCAST_NS_TIMESTAMP_TO_US_ON_WRITE, TableProperties |
2796 | | - |
| 2794 | +def _resolve_write_setup(table_metadata: TableMetadata) -> tuple[FileFormat, FileFormatModel, LocationProvider, Schema]: |
2797 | 2795 | file_format = FileFormat( |
2798 | | - table_metadata.properties.get( |
2799 | | - TableProperties.WRITE_FILE_FORMAT, |
2800 | | - TableProperties.WRITE_FILE_FORMAT_DEFAULT, |
2801 | | - ) |
| 2796 | + table_metadata.properties.get(TableProperties.WRITE_FILE_FORMAT, TableProperties.WRITE_FILE_FORMAT_DEFAULT) |
2802 | 2797 | ) |
2803 | 2798 | format_model = FileFormatFactory.get(file_format) |
2804 | 2799 | location_provider = load_location_provider(table_location=table_metadata.location, table_properties=table_metadata.properties) |
| 2800 | + table_schema = table_metadata.schema() |
| 2801 | + if (sanitized_schema := sanitize_column_names(table_schema)) != table_schema: |
| 2802 | + file_schema = sanitized_schema |
| 2803 | + else: |
| 2804 | + file_schema = table_schema |
| 2805 | + return file_format, format_model, location_provider, file_schema |
2805 | 2806 |
|
2806 | | - def write_data_file(task: WriteTask) -> DataFile: |
2807 | | - table_schema = table_metadata.schema() |
2808 | | - if (sanitized_schema := sanitize_column_names(table_schema)) != table_schema: |
2809 | | - file_schema = sanitized_schema |
2810 | | - else: |
2811 | | - file_schema = table_schema |
2812 | 2807 |
|
| 2808 | +def write_file(io: FileIO, table_metadata: TableMetadata, tasks: Iterator[WriteTask]) -> Iterator[DataFile]: |
| 2809 | + file_format, format_model, location_provider, file_schema = _resolve_write_setup(table_metadata) |
| 2810 | + |
| 2811 | + def write_data_file(task: WriteTask) -> DataFile: |
2813 | 2812 | downcast_ns_timestamp_to_us = Config().get_bool(DOWNCAST_NS_TIMESTAMP_TO_US_ON_WRITE) or False |
2814 | 2813 | batches = [ |
2815 | 2814 | _to_requested_schema( |
@@ -3070,18 +3069,7 @@ def _dataframe_to_data_files( |
3070 | 3069 | "Materialise the reader as a pa.Table first, or follow " |
3071 | 3070 | "https://github.com/apache/iceberg-python/issues/2152 for partitioned streaming support." |
3072 | 3071 | ) |
3073 | | - file_format = FileFormat( |
3074 | | - table_metadata.properties.get(TableProperties.WRITE_FILE_FORMAT, TableProperties.WRITE_FILE_FORMAT_DEFAULT) |
3075 | | - ) |
3076 | | - format_model = FileFormatFactory.get(file_format) |
3077 | | - location_provider = load_location_provider( |
3078 | | - table_location=table_metadata.location, table_properties=table_metadata.properties |
3079 | | - ) |
3080 | | - table_schema = table_metadata.schema() |
3081 | | - if (sanitized_schema := sanitize_column_names(table_schema)) != table_schema: |
3082 | | - file_schema = sanitized_schema |
3083 | | - else: |
3084 | | - file_schema = table_schema |
| 3072 | + file_format, format_model, location_provider, file_schema = _resolve_write_setup(table_metadata) |
3085 | 3073 |
|
3086 | 3074 | batches = iter(df) |
3087 | 3075 | for batch in batches: |
|
0 commit comments