Last active
June 29, 2026 19:51
-
-
Save u8sand/44d39f35c779192f4a34bf5279356ae3 to your computer and use it in GitHub Desktop.
A micro-library for quick pandas dataframe => postgres upload
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| ''' | |
| This snippet can be used for quickly copying a pandas dataframe | |
| into postgres via unix pipe. I expect this to be faster than other | |
| approaches since the writing can happen entirely in native code: | |
| - os.pipe is facilitated by the kernel | |
| - psycopg2 copy_from uses libpg | |
| - np.savetxt uses libnumpy | |
| Usage: | |
| ```python | |
| import psycopg2 | |
| import pandas as pd | |
| import numpy as np | |
| # or just paste the functions here | |
| from df2pg import copy_from_df | |
| # some random postgres connection | |
| con = psycopg2.connect('postgresql://postgres:postgres@localhost:5432/postgres') | |
| # some random schema | |
| cur = con.cursor() | |
| cur.execute(""" | |
| create table "data" ( | |
| "x" decimal, | |
| "y" decimal, | |
| "z" decimal | |
| ); | |
| """) | |
| con.commit() | |
| # some random dataframe | |
| df = pd.DataFrame(np.random.normal(size=(1000000, 3)), columns=['x', 'y', 'z'], copy=False) | |
| # actually copy the data | |
| copy_from_df(con, 'data', df) | |
| # verify it worked | |
| cur = con.cursor() | |
| cur.execute('select * from data limit 5') | |
| cur.fetchall() | |
| df.head() | |
| ''' | |
| import uuid | |
| import typing as t | |
| import psycopg2 | |
| FileDescriptor = t.Union[int, str] | |
| class OnConflictSpec(t.TypedDict): | |
| # columns that conflict | |
| conflict: t.Tuple[str] | |
| # Tuple to update specific columns | |
| # True to update all non-conflicting columns | |
| # False to ignore (default) | |
| update: t.Optional[t.Union[t.Tuple[str], bool]] | |
| def peek(it: t.Iterator): | |
| import itertools | |
| first = next(it) | |
| return first, itertools.chain((first,), it) | |
| def copy_from_tsv( | |
| con: 'psycopg2.connection', | |
| table: str, | |
| r: FileDescriptor, | |
| columns: t.Optional[list[str]] = None, | |
| on: t.Optional[OnConflictSpec] = None, | |
| setup: t.Optional[str] = None, | |
| ): | |
| ''' Copy from a file descriptor into a postgres database table through as psycopg2 connection object | |
| :param con: The psycopg2.connect object | |
| :param table: The table top copy into | |
| :param r: An file descriptor to be opened in read mode | |
| :param columns: The columns being copied | |
| :param on: What to do on conflict (see OnConflictSpec) | |
| :param setup: Commands to run in the cursor before ingest | |
| ''' | |
| import os | |
| with con.cursor() as cur: | |
| if setup: cur.execute(setup) | |
| with os.fdopen(r, 'rb', buffering=0, closefd=True) as fr: | |
| tsv_columns = list(map(bytes.decode, fr.readline().strip().split(b'\t'))) | |
| if columns is None: | |
| columns = tsv_columns | |
| else: | |
| columns = [c for c in tsv_columns if c in columns] | |
| if on: | |
| on_update = on.get('update') | |
| if on_update is True: | |
| on_update = {c for c in columns if c not in on['conflict']} | |
| tmp_table = f"tmp_{str(uuid.uuid4()).partition('-')[0]}" | |
| cur.copy_expert( | |
| sql=f''' | |
| CREATE TEMPORARY TABLE {tmp_table} as table {table} WITH NO DATA; | |
| COPY {tmp_table} ({",".join(f'"{c}"' for c in columns)}) | |
| FROM STDIN WITH CSV DELIMITER E'\\t'; | |
| INSERT INTO {table} ({",".join(f'"{c}"' for c in columns)}) | |
| SELECT {",".join(f'"{c}"' for c in columns)} | |
| FROM {tmp_table} | |
| ON CONFLICT ({",".join(f'"{c}"' for c in on["conflict"])}) | |
| DO { | |
| f"""UPDATE SET {", ".join(f'"{c}" = EXCLUDED."{c}"' for c in columns if c in on_update)}""" | |
| if on_update else f"""NOTHING"""}; | |
| DROP TABLE {tmp_table}; | |
| ''', | |
| file=fr, | |
| ) | |
| else: | |
| cur.copy_expert( | |
| sql=f''' | |
| COPY {table} ({",".join(f'"{c}"' for c in columns)}) | |
| FROM STDIN WITH CSV DELIMITER E'\\t' | |
| ''', | |
| file=fr, | |
| ) | |
| con.commit() | |
| def copy_from_records( | |
| con: 'psycopg2.connection', | |
| table: str, | |
| records: t.Iterable[dict], | |
| columns: t.Optional[list[str]] = None, | |
| on: t.Optional[OnConflictSpec] = None, | |
| setup: t.Optional[str] = None, | |
| ): | |
| ''' Copy from records into a postgres database table through as psycopg2 connection object. | |
| This is done by constructing a unix pipe, writing the records with csv writer | |
| into the pipe while loading from the pipe into postgres at the same time. | |
| :param con: The psycopg2.connect object | |
| :param table: The table to write the pandas dataframe into | |
| :param records: An iterable of records to write | |
| :param columns: The columns being written into the table | |
| :param on: What to do on conflict (see OnConflictSpec) | |
| :param setup: Commands to run in the cursor before ingest | |
| ''' | |
| import os, csv, threading | |
| records_it = iter(records) | |
| if columns is None: | |
| try: | |
| first, records_it = peek(records_it) | |
| columns = tuple(first.keys()) | |
| except StopIteration: | |
| return | |
| # | |
| r, w = os.pipe() | |
| # we copy_from_tsv with the read end of this pipe in | |
| # another thread | |
| rt = threading.Thread( | |
| target=copy_from_tsv, | |
| args=(con, table, r, None, on, setup), | |
| ) | |
| rt.start() | |
| try: | |
| # we write to the write end of this pipe in this thread | |
| with os.fdopen(w, 'w', closefd=True) as fw: | |
| writer = csv.DictWriter(fw, fieldnames=columns, delimiter='\t') | |
| writer.writeheader() | |
| writer.writerows(records_it) | |
| finally: | |
| # we wait for the copy_from_tsv thread to finish | |
| rt.join() | |
| try: | |
| import pandas as pd | |
| except ImportError: | |
| def copy_from_df(*args, **kwargs): | |
| raise RuntimeError('pandas is required') | |
| else: | |
| def copy_from_df( | |
| con: 'psycopg2.connection', | |
| table: str, | |
| df: pd.DataFrame, | |
| float_format: str = '%g', | |
| native: bool = True, | |
| on: t.Optional[OnConflictSpec] = None, | |
| setup: t.Optional[str] = None, | |
| ): | |
| ''' Copy from a pandas dataframe into a postgres database table through as psycopg2 connection object. | |
| This is done by constructing a unix pipe, writing the data frame | |
| into the pipe while loading from the pipe into postgres at the same time. | |
| :param con: The psycopg2.connect object | |
| :param table: The table to write the pandas dataframe into | |
| :param df: The pandas dataframe to write to the database | |
| :param native: Use native write, this should be preferred as it performs writes | |
| with numpy (in C), but it might be worse if you have python | |
| objects in your data frame. | |
| :param on: What to do on conflict (see OnConflictSpec) | |
| :param setup: Commands to run in the cursor before ingest | |
| ''' | |
| import os, numpy as np, threading | |
| r, w = os.pipe() | |
| # we copy_from_tsv with the read end of this pipe in | |
| # another thread | |
| rt = threading.Thread( | |
| target=copy_from_tsv, | |
| args=(con, table, r, None, on, setup), | |
| ) | |
| rt.start() | |
| try: | |
| # we write to the write end of this pipe in this thread | |
| with os.fdopen(w, 'wb', buffering=0, closefd=True) as fw: | |
| if native: | |
| # np.savetxt is typically faster than pd.to_csv as more | |
| # of it seems to be implemented in native code | |
| np.savetxt( | |
| fw, df.values, | |
| fmt=float_format, | |
| delimiter='\t', | |
| newline='\n', | |
| comments='', | |
| header='\t'.join(df.columns.tolist()), | |
| ) | |
| else: | |
| # pandas to_csv is probably more compatible with a wider | |
| # range of data types but is slower | |
| df.to_csv( | |
| fw, | |
| float_format=float_format, | |
| sep='\t', | |
| index=None, | |
| ) | |
| finally: | |
| # we wait for the copy_from_tsv thread to finish | |
| rt.join() |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| # Allow you to pip install this gist with: | |
| # pip install git+https://gist.github.com/44d39f35c779192f4a34bf5279356ae3.git | |
| [project] | |
| name = 'df2pg' | |
| version = '1.1.4' | |
| description = '' | |
| authors = [ | |
| {name = "Daniel J. B. Clarke", email = "danieljbclarkemssm@gmail.com"} | |
| ] | |
| requires-python = ">=3.9" | |
| dependencies = [ | |
| # one of these is required, not sure how to specify this at the moment | |
| # "psycopg2", | |
| # "psycopg2-binary", | |
| "pandas ; extra == 'copy_from_df'", | |
| "numpy ; extra == 'copy_from_df'" | |
| ] | |
| [project.optional-dependencies] | |
| copy_from_df = [ "pandas", "numpy" ] | |
| [build-system] | |
| requires = ["poetry-core>=2.0.0,<3.0.0"] | |
| build-backend = "poetry.core.masonry.api" |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment