Skip to content

bq dataframe - #1630

Merged
fiedlerNr9 merged 2 commits into
mainfrom
jan/bq-dataframe
Sep 30, 2026
Merged

fiedlerNr9 merged 2 commits into
mainfrom
jan/bq-dataframe

Conversation

@fiedlerNr9

@fiedlerNr9 fiedlerNr9 commented Sep 30, 2026 •

Copy link
Copy Markdown
Contributor

Summary

Adds the missing BigQuery-to-pandas DataFrame decoder.

BigQuery connector outputs use bq://project:dataset.table URIs, which previously fell back to the Parquet decoder and failed with Protocol not known: bq.

The decoder:

  • Reads results through the BigQuery Storage API
  • Supports column projection
  • Handles empty result tables
  • Registers automatically with the bq protocol

Testing

  • Added decoder, URI parsing, projection, and empty-table tests
  • All 42 BigQuery plugin tests pass

Before

Running this:

import flyte
from flyte.io import DataFrame
from flyteplugins.bigquery import BigQueryConfig, BigQueryTask
import pandas as pd

env = flyte.TaskEnvironment(
    "bq-test",
    image=flyte.Image.from_debian_base().with_pip_packages(
        "pandas", "pyarrow", "flyteplugins-bigquery"
    ),
)


config = BigQueryConfig(
    ProjectID="dogfood-gcp-dataplane",
    Location="US",
)

query_task = BigQueryTask(
    name="jan-query",
    query_template="SELECT * FROM `dogfood-gcp-dataplane.dataset.demo_table` LIMIT 1000",
    plugin_config=config,
    output_dataframe_type=DataFrame,
)


@env.task
async def main() -> pd.DataFrame:
    df_bq = await query_task()
    df = await df_bq.open(pd.DataFrame).all()
    print(df)
    return df


if __name__ == "__main__":
    flyte.init_from_config("dogfood-gcp.yaml")
    run = flyte.run(main)
    print(run.url)

-> Failed here

After

Running this

import flyte
from flyte.io import DataFrame
from flyteplugins.bigquery import BigQueryConfig, BigQueryTask
import pandas as pd

env = flyte.TaskEnvironment(
    "bq-test",
    image=flyte.Image.from_debian_base()
    .with_apt_packages("git")
    .with_pip_packages(
        "pandas",
        "pyarrow",
        "flyteplugins-bigquery @ git+https://github.com/flyteorg/flyte-sdk.git@jan/bq-dataframe#subdirectory=plugins/bigquery",
    ),
)


config = BigQueryConfig(
    ProjectID="dogfood-gcp-dataplane",
    Location="US",
)

query_task = BigQueryTask(
    name="jan-query",
    query_template="SELECT * FROM `dogfood-gcp-dataplane.dataset.demo_table` LIMIT 1000",
    plugin_config=config,
    output_dataframe_type=pd.DataFrame,
)


@env.task
async def main() -> pd.DataFrame:
    df_bq = await query_task()
    print(df_bq)
    df = df_bq
    # df = await df_bq.open(pd.DataFrame).all()
    print(df)
    return df


if __name__ == "__main__":
    flyte.init_from_config("dogfood-gcp.yaml")
    run = flyte.run(main)
    print(run.url)

-> Succeeded here

Signed-off-by: fiedlerNr9 <jan@union.ai>
Signed-off-by: fiedlerNr9 <jan@union.ai>
@fiedlerNr9
fiedlerNr9 merged commit 060d604 into main Sep 30, 2026
67 checks passed
@fiedlerNr9
fiedlerNr9 deleted the jan/bq-dataframe branch September 30, 2026 13:45
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants