Skip to content
Merged
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
8 changes: 5 additions & 3 deletions sdks/python/apache_beam/io/parquetio.py
Original file line number Diff line number Diff line change
Expand Up @@ -52,11 +52,13 @@

try:
import pyarrow as pa
paTable = pa.Table
import pyarrow.parquet as pq
# pylint: disable=ungrouped-imports
from apache_beam.typehints import arrow_type_compatibility
except ImportError:
pa = None
paTable = None
pq = None
ARROW_MAJOR_VERSION = None
arrow_type_compatibility = None
Expand Down Expand Up @@ -176,7 +178,7 @@ def __init__(self, beam_type):
self._beam_type = beam_type

@DoFn.yields_batches
def process(self, element) -> Iterator[pa.Table]:
def process(self, element) -> Iterator[paTable]:
yield element

def infer_output_type(self, input_type):
Expand All @@ -185,7 +187,7 @@ def infer_output_type(self, input_type):

class _BeamRowsToArrowTable(DoFn):
@DoFn.yields_elements
def process_batch(self, element: pa.Table) -> Iterator[pa.Table]:
def process_batch(self, element: paTable) -> Iterator[paTable]:
yield element


Expand Down Expand Up @@ -845,7 +847,7 @@ def open(self, temp_path):
use_deprecated_int96_timestamps=self._use_deprecated_int96_timestamps,
use_compliant_nested_type=self._use_compliant_nested_type)

def write_record(self, writer, table: pa.Table):
def write_record(self, writer, table: paTable):
writer.write_table(table)

def close(self, writer):
Expand Down
Loading