Pandas Dataframe to Apache Beam PCollection conversion problem

Viewed 1935

I'm trying to convert a pandas DataFrame to a PCollection from Apache Beam. Unfortunately, when I use to_pcollection() function, I get the following error:

AttributeError: 'DataFrame' object has no attribute '_expr'

Does anyone know how to solve it? I'm using pandas=1.1.4, beam=2.25.0 and Python 3.6.9.

2 Answers

to_pcollection was only ever intended to apply to Beam's deferred Dataframes, but looking at this it makes sense that it should work, and isn't obvious how to do manually. https://github.com/apache/beam/pull/14170 should fix this.

I get this problem when I use a "native" Pandas dataframe instead of a dataframe created by to_dataframe within Beam. I suspect that the dataframe created by Beam wraps or subclasses a Pandas dataframe with new attributes (like _expr) that the native Pandas dataframe doesn't have.

The real answer involves knowing how to use apache_beam.dataframe.convert.to_dataframe, but I can't figure out how to set the proxy object correctly (I get Singleton errors when I try to later use to_pcollection). So since I can't get the "right" way to to work in 2.25.0 (I'm new to Beam and Pandas—and don't know how proxy objects work—so take all this with a grain of salt), I use this workaround:

class SomeDoFn(beam.DoFn):
    def process(self, pair): # pair is a key/value tuple
        df = pd.DataFrame(pair[1]) # just the array of values

        ## do something with the dataframe
        ...

        records = df.to_dict('records')

        # return a tuple with the same shape as the one we received
        return [(rec["key"], rec) for rec in records]

which I invoke with something like this:

rows = (
    pcoll
    | beam.ParDo(SomeDoFn())
)

I hope others will give you a better answer than this workaround.

Related