How can I get the following SQLAlchemy expression to work with bigquery (similarly to how it works with SQLite):
conn.execute(
data.update().
where(sqlalchemy.tuple_(data.c.person_id, data.c.person_name).
not_in(sqlalchemy.sql.select(people.c.id, people.c.name))).
values(invalid=True))
This results in the error Subquery of type IN must have only one output column when run with BigQuery:
sqlalchemy.exc.DatabaseError: (google.cloud.bigquery.dbapi.exceptions.DatabaseError) 400 Subquery of type IN must have only one output column at [1:98]
Location: US
Job ID: 4331d568-e06f-41fa-813b-9754574587d7
[SQL: UPDATE `data` SET `invalid`=%(invalid:BOOL)s WHERE ((`data`.`person_id`, `data`.`person_name`) NOT IN (SELECT `people`.`id`, `people`.`name`
FROM `people`))]
[parameters: {'invalid': True}]
(Background on this error at: https://sqlalche.me/e/14/4xp6)
Using SELECT AS STRUCT in the subquery resolves this error -- how can I do this with SQLAlchemy (and why doesn't SQLAlchemy automatically do this)? For example this works in BigQuery:
UPDATE `jdimatteo-v.scratch.data`
SET `invalid`=True
WHERE (person_id, person_name)
NOT IN (SELECT AS STRUCT id, name
FROM `jdimatteo-v.scratch.people`
)
Here is a full working example with SQLite, with the connection to BigQuery that results in the above error commented out:
#!/usr/bin/env python3
import sqlalchemy # pip install SQLAlchemy==1.4.27
engine = sqlalchemy.create_engine('sqlite://')
# engine = sqlalchemy.create_engine('bigquery://jdimatteo-v/scratch') # pip install sqlalchemy-bigquery==1.4.3
metadata_obj = sqlalchemy.MetaData()
people = sqlalchemy.Table(
'people', metadata_obj,
sqlalchemy.Column('id', sqlalchemy.Integer),
sqlalchemy.Column('name', sqlalchemy.String),
)
data = sqlalchemy.Table(
'data', metadata_obj,
sqlalchemy.Column('person_id', sqlalchemy.Integer),
sqlalchemy.Column('person_name', sqlalchemy.String),
sqlalchemy.Column('data_foo', sqlalchemy.String),
sqlalchemy.Column('invalid', sqlalchemy.Boolean),
)
metadata_obj.create_all(engine)
conn = engine.connect()
def create_records():
conn.execute(people.delete().where(True == True))
conn.execute(people.insert().values(id=1, name='Mary'))
conn.execute(people.insert().values(id=2, name='James'))
conn.execute(data.delete().where(True == True))
conn.execute(data.insert().values(person_id=1, person_name='Mary', data_foo='good foo', invalid=None))
conn.execute(data.insert().values(person_id=42, person_name='Bob', data_foo='chop suey', invalid=None))
conn.execute(data.insert().values(person_id=1, person_name='James', data_foo='mixed up', invalid=None))
def dynamic_update():
conn.execute(
data.update().
where(sqlalchemy.tuple_(data.c.person_id, data.c.person_name).
not_in(sqlalchemy.sql.select(people.c.id, people.c.name))).
values(invalid=True))
def print_records(msg):
print(f'{msg}:')
print(" people:")
for person in conn.execute(sqlalchemy.sql.select(people)):
print(' ', person)
print(" data:")
for datum in conn.execute(sqlalchemy.sql.select(data)):
print(' ', datum)
print()
create_records()
print_records('initial values')
dynamic_update()
print_records('after update')