hussh
Agent One
Products
PuppyTagShop
Marketplace
OverviewWhite PagesYellow Pages
For Business
For Advisors & RIAFor BrandsFor Agent BuildersFor DefenseOur PartnersPartner with Hussh
Blog
Guide
A-ZPCHPYour informationBuilding in the openConsent & Privacy
Company
Our StoryTeamCareersPressMediaContact
Get Agent One
Agent One
Products
PuppyTagShop
Marketplace
OverviewWhite PagesYellow Pages
For Business
For Advisors & RIAFor BrandsFor Agent BuildersFor DefenseOur PartnersPartner with Hussh
Blog
Guide
A-ZPCHPYour informationBuilding in the openConsent & Privacy
Company
Our StoryTeamCareersPressMediaContact
The magazine
OpenAIAutomationFastapi

How to Set Up an Automated Batch Processing Pipeline for the OpenAI API (GPT-4o-mini)

A complete guide to an automated batch processing pipeline with the OpenAI Batch API, Supabase and FastAPI: from creating batch jobs to saving the results.

Omkar Malpure·November 14, 2024·9 min read
How to Set Up an Automated Batch Processing Pipeline for the OpenAI API (GPT-4o-mini)

Introduction

Now a days utilising LLM’s for processing data on enterprise level has become a topic of dicussion . Lot of companies are heavliy investing in this segment of LLM finetuning , pre training , so that they are not left behind in the race .

But that being said LLM’s require a significant amount of investment , and even the cost for hosting an LLM or for inferencing is quite high due to the huge compute resource requirements of the LLM models .

In this blog we are going to explore , how we can process huge volume of data utilsing OpenAI API , we will be utilising supabase database to store the data . Supabase is a relational database .

If you are new to supabase , I will reccommend to get familiar with supabase , so that you can follow this tutorial with ease .

That being said , I am just using supabase here as an example , If you are experienced with utilising any kind of database may it be relational or non - relational databases , like AWS RDS (Relational Database service) , Dynamo DB (NoSql or Non-relational database) . You just have to implement the ingestion layer corresponding to your database .

I will be utilising python programming language and FastAPI framework , to build API’s , to ingest the data from the database , process the data to bring into the format required by OpenAI batch processing service .

I would also be utilising concepts such as triggers and database functions , and these would be written in PostgreSql , again if you are not familiar with it , no need to worry , I have added comments for each line which can be helpful to understand .

It would be useful if you can go through basics of SQL and get an understanding of SELECT , CREATE , FOR EACH , DELETE , etc .

You can go through resources like w3 schools or tutorials points or there are plenty of resources available on youtube as well .

Why use batch processing?

Imagine you have 50,000 data points to process with OpenAI’s GPT-4o-mini API. The first idea that comes to mind is a simple loop that sends one request at a time. That approach has a few drawbacks:

  • If your code is deployed on a server, it keeps executing for all 50,000 requests, and that many data points means a lot of tokens to process, which can add up to a hefty bill. You can refer to the pricing page OpenAI publishes.
  • If the code throws an error partway through, further execution stops. You might start 50,000 requests at night, hoping they finish by morning, and wake up to find only 5,000 completed.
  • If you run it on your local machine, any internet failure stops the execution.

OpenAI’s Batch API avoids all of this. Batch requests cost 50% less than the same requests sent one by one, and they run against a separate pool of higher rate limits. To follow along, set up an OpenAI account and create an API key at https://platform.openai.com/settings/organization/api-keys.

How a batch request is shaped

The batch file, in the jsonl format, should contain one line (json object) per request. Each request is defined as such:

{
 “custom_id”: <REQUEST_ID>,
 “method”: “POST”,
 “url”: “/v1/chat/completions”,
 “body”: {
 “model”: 'gpt-4o-mini',
  "temperature": 0.2,
 “messages”: [
    {
       "role": "user",
       "content": f"{text}"
    },
  ]
 }
}

Note: the request ID should be unique per batch. This is what you can use to match results to the initial input files, as requests will not be returned in the same order.

You don’t have to save this file to disk. Building it in memory is what lets the whole pipeline run on a server and be automated with triggers and cron jobs, as you will see below.

Usecase : In this example, we will use gpt-4o-mini to extract movie categories from a description of the movie. We will also extract a 1-sentence summary from this description.

Now let’s see into how we can create the automated pipeline .

Step 1: Set up database functions and triggers in Supabase

On the left side you will find an option to select SQL editor .

Supabase

Once you are in the editor create a new snippet as you can in the above image .

Let’s first create the Batch processing function .

The purpose of this function would be to ingest the 1000(as supabase allows minimum 1000 rows in its fetch operation but it could be changed if you want it for enterprise you can check supabase documentation) rows of data , and pass this data to the API , an endpoint that we would create using FastApi .

Here is the PostgreSQL snippet

CREATE OR REPLACE FUNCTION process_batch_data_v3()
RETURNS json AS $$
DECLARE
    last_processed_id BIGINT;
    current_batch_data json;
    start_id BIGINT;
    end_id BIGINT;
BEGIN
    WITH batch_rows AS (
        SELECT 
            imdb_id, -- row names
            "Description",
            ROW_NUMBER() OVER (ORDER BY imdb_id) as rn
        FROM your_table_name
        LIMIT 1000
    )
    SELECT 
        json_build_object(
            'data', (SELECT array_to_json(array_agg(row_to_json(batch_rows)))) ,
            'min_id', MIN(imdb_id),
            'max_id', MAX(imdb_id)
        )
    INTO current_batch_data
    FROM batch_rows;

    PERFORM http_post(
        'Our FastApi endpoint',
        (current_batch_data)::text,
        'application/json'
    );
    
    DELETE FROM your_table_name
    WHERE id BETWEEN start_id AND end_id;

    RETURN current_batch_data;
END;
$$ LANGUAGE plpgsql;

Replace your_table_name with your actual table name that you have give for your table.

Just execute the above code and a data ingestion database function will be created for 1000 rows , which can be passed for batching using OpenAI API.

Once the rows are passed for batch processing those rows would also be deleted , for taking up next 1000 rows.(Another approach would be to add another column to add a check that batch processing for this row is completed.)

Now lets setup the database function for the trigger which would tell us when 1000 rows are available in our database .

Database function for Trigger

CREATE OR REPLACE FUNCTION check_row_count()
RETURNS TRIGGER AS $$
DECLARE
    row_count integer;
BEGIN
    SELECT COUNT(*) INTO row_count FROM your_table_name;
    
    IF row_count = 1000 THEN
        PERFORM process_batch_data_v3();
    END IF;
    
    RETURN NEW;
END;
$$ LANGUAGE plpgsql;

Replace your_table_name with your actual table name .

Just execute the above code and your database function which would act as a trigger would be created .

Let’s setup the trigger for the above database function check_row_count().

CREATE TRIGGER check_row_count_trigger
AFTER INSERT ON imdb_dataset
FOR EACH ROW
EXECUTE FUNCTION check_row_count();

Step 2: Create the FastAPI endpoint that creates the batch job

You can run this endpoint locally or you can deploy this endpoint on hugging face spaces on docker and utilise it as a server .

I will give you the code for setting it up locally .

Lets first install all the libraries needed.

pip install supabase openai fastapi python-dotenv uvicorn

Here is the code below

from fastapi import FastAPI, Request
import os
import json
import io
from dotenv import load_dotenv
from supabase import create_client, Client
from openai import Client as OpenAIClient
import logging

app = FastAPI()

# Load environment variables
load_dotenv()

# Initialize OpenAI client
client = OpenAIClient(api_key=os.getenv('OPENAI_API_KEY'), organization=os.getenv('ORG_ID'))

# Initialize Supabase client
url: str = os.getenv('SUPABASE_URL')
key: str = os.getenv('SUPABASE_KEY')
supabase: Client = create_client(url, key)

@app.post("/test/v1")
async def testv1(request: Request):
    system_prompt = '''
        Your goal is to extract movie categories from movie descriptions, as well as a 1-sentence summary for these movies.
        You will be provided with a movie description, and you will output a json object containing the following information:
        
        {
            categories: string[] // Array of categories based on the movie description,
            summary: string // 1-sentence summary of the movie based on the movie description
        }
        
        Categories refer to the genre or type of the movie, like "action", "romance", "comedy", etc. Keep category names simple and use only lower case letters.
        Movies can have several categories, but try to keep it under 3-4. Only mention the categories that are the most obvious based on the description.
    '''
    
    dataset = await request.json()
    tasks = []
    for ds in dataset.get('data', []):
        imdb_id = ds.get('imdb_id')
        description = ds.get('Description')
        task = {
            "custom_id": f"task-{imdb_id}",
            "method": "POST",
            "url": "/v1/chat/completions",
            "body": {
                "model": "gpt-4o-mini",
                "temperature": 0.1,
                "response_format": { 
                    "type": "json_object"
                },
                "messages": [
                    {
                        "role": "system",
                        "content": system_prompt
                    },
                    {
                        "role": "user",
                        "content": description
                    }
                ],
            }
        }
        tasks.append(task)
    
    # Write tasks to a JSON object
    json_obj = io.BytesIO()
    for obj in tasks:
        json_obj.write((json.dumps(obj) + '\n').encode('utf-8'))

    # Create a batch job
    batch_file = client.files.create(
        file=json_obj,
        purpose="batch"
    )
    batch_job = client.batches.create(
        input_file_id=batch_file.id,
        endpoint="/v1/chat/completions",
        completion_window="24h"
    )

    save_data = {
        'batch_job_id': batch_job.id,
        "batch_job_status": False
    }

    # Save batch job details to Supabase
    response = supabase.table("batch_processing_details").insert(save_data).execute()
        
    return {'data': 'Batch job is scheduled!'}

if __name__ == "__main__":
    import uvicorn
    uvicorn.run(app, host="0.0.0.0", port=8000)

So above is the code utilising fastapi and openai batch processing .

I will explain what exactly the code does , so the below part of the code

dataset = await request.json() 

extracts the data from the sent from the database function.

tasks = []
    for ds in dataset.get('data', []):
        imdb_id = ds.get('imdb_id')
        description = ds.get('Description')
        task = {
            "custom_id": f"task-{imdb_id}",
            "method": "POST",
            "url": "/v1/chat/completions",
            "body": {
                "model": "gpt-4o-mini",
                "temperature": 0.1,
                "response_format": { 
                    "type": "json_object"
                },
                "messages": [
                    {
                        "role": "system",
                        "content": system_prompt
                    },
                    {
                        "role": "user",
                        "content": description
                    }
                ],
            }
        }
        tasks.append(task)

This code goes through each row and constructs a openai request json which is a required format for openai.

# Write tasks to a JSON object
    json_obj = io.BytesIO()
    for obj in tasks:
        json_obj.write((json.dumps(obj) + '\n').encode('utf-8'))

Above code loads json required to pass for batch process in memory without requiring to save it locally or even on a cloud bucket . But this depends on the amount of data you would ingest and as long as you have enough memory on server as well .

# Create a batch job
    batch_file = client.files.create(
        file=json_obj,
        purpose="batch"
    )
    batch_job = client.batches.create(
        input_file_id=batch_file.id,
        endpoint="/v1/chat/completions",
        completion_window="24h"
    )

    save_data = {
        'batch_job_id': batch_job.id,
        "batch_job_status": False
    }

    # Save batch job details to Supabase
    response = supabase.table("batch_processing_details").insert(save_data).execute()

Above code is a standard code for creating a batch job where we pass the in-memory json object we have created , and then we create the batch by client.batches.create with , completion endpoint and completion_window parameters .

Now we extract the batch job id , and we will save this data onto a seperate table batch_processing_details , that you will have create along with columns batch_job_id , batch_job_status , id , created_at .

You can run the above code using below command

python your_file_name.py

you will get the endpoint below one you run the command locally .

Supabase

Now you have to pass the endpoint I have given below to your database function . Here is the endpoint .

http://0.0.0.0:8000/test/v1

You have to pass the endpoint in this database function .

CREATE OR REPLACE FUNCTION process_batch_data_v3()
RETURNS json AS $$
DECLARE
    last_processed_id BIGINT;
    current_batch_data json;
    start_id BIGINT;
    end_id BIGINT;
BEGIN

    -- Get the batch data
    WITH batch_rows AS (
        SELECT 
            imdb_id,
            "Description",
            ROW_NUMBER() OVER (ORDER BY imdb_id) as rn
        FROM imdb_dataset
        LIMIT 1000
    )
    SELECT 
        json_build_object(
            'data', (SELECT array_to_json(array_agg(row_to_json(batch_rows)))) ,
            'min_id', MIN(imdb_id),
            'max_id', MAX(imdb_id)
        )
    INTO current_batch_data
    FROM batch_rows;

    PERFORM http_post(
-- pass the endpoint here as given
        'http://0.0.0.0:8000/test/v1',
        (current_batch_data)::text,
        'application/json'
    );
    
    DELETE FROM your_table
    WHERE id BETWEEN start_id AND end_id;

    -- Return the batch data
    RETURN current_batch_data;
END;
$$ LANGUAGE plpgsql;

Replace your_table_name with your actual table name .

Run this code again on the SQL editor tp update the function with desired endpoint .

Step 3: Load the IMDB data into Supabase

Note: If you already have a different database filled with data, then you can just skip this step.

Let’s dowload the data and store it into supabase .

You can download your data here: Dataset

You may have to modify the dataset as it contains one empty column when you download it .

If you are new to supabase then you may have to create a new project and add a new table . You can just go through this 2 mins tutorial : tutorial

Above tutorial will show you to setup the project and create tables .

Now create a table with these columns imdb_id(as primary key) , created_at (Both these columns would be already present just rename id as imdb_id) , Title , Certificate , Duration , Genre , Rate(datatype float8) , Metascore(int8) , Description(text) , Cast(text) , Info(text) .

Now upload your csv data that you had downloaded .

Go to Insert>Import Data from CSV , as shown below.

Supabase

Now upload your data , by clicking on browse .

Once the whole 1000 rows are uploaded , the database function would be triggered and the 1000 rows data would be passed to the endpoint , and they would be sent for batch processing .

Step 4: Create a table for the batch outputs

Next, set up the table where you will save the outputs that come back from batch processing (OpenAI API, GPT-4o-mini). Its shape depends on your data if you are using a totally different dataset.

For the IMDB data from Step 3, create a table with the columns id, created_at, description, categories and summary. Set the id column without Is Identity, as shown below.

Supabase

Step 5: Create the database function that checks for completed batches

This database function reads each batch_job_id from the batch_processing_details table you created in Step 2 and asks an API endpoint whether that batch job has finished.

In the SQL editor run this code and your database function would be created.

CREATE OR REPLACE FUNCTION check_batch_processing_completion()
RETURNS void AS $$
DECLARE
    row_record RECORD;
    api_response JSON;
    row_json JSON;
    batch_status text;
BEGIN
    FOR row_record IN SELECT * FROM your_table_name WHERE batch_job_status = FALSE
    LOOP
        -- Make an API call to the external endpoint
        -- Replace 'your_api_endpoint' with the actual API endpoint
        -- and 'column_to_send' with the column name to send to the API
        row_json := json_build_object('batch_job_id', row_record.batch_job_id);
        SELECT content::json INTO api_response
        FROM http_post(
            'http://0.0.0.0:8000/test/batch_processing_result',
            row_json::text,
            'application/json'
        );

        batch_status := api_response->>'batch_job_status';
        -- Update the column_to_check with the API response
        UPDATE your_table_name
        SET batch_job_status = CASE 
            WHEN batch_status = 'completed' THEN TRUE
            WHEN batch_status = 'notcompleted' THEN FALSE
            ELSE batch_job_status  -- Keep existing value if response is invalid
        END
        WHERE id = row_record.id;
    END LOOP;
END;
$$ LANGUAGE plpgsql;

Replace your_table_name with the name of your batch_processing_details table. The function goes through every batch job that is not marked complete, passes its batch id to the API from Step 7, and marks the row complete once the API reports that the batch job has finished. The API then pushes the results into the outputs table.

Step 6: Schedule the cron job

SELECT cron.schedule('0 */2 * * *', 'SELECT check_batch_processing_completion();');

Now this cron job is setup which will run the database function every 2 hours.

If you would rather poll from Python than use a cron job, you can check the status of a batch job in a loop:

batch_job = client.batches.retrieve(id)

while batch_job.status == 'in_progress':
    batch_job = client.batches.retrieve(id)
    print(batch_job.status)
    if batch_job.status == 'completed':
        break
    time.sleep(60)

Step 7: Set up the FastAPI endpoint that saves the outputs

from fastapi import FastAPI, Request, BackgroundTasks
from supabase import create_client, Client
from openai import Client as OpenAIClient
import json
import uvicorn
from typing import Dict, List, Optional
import os
from dotenv import load_dotenv

# Load environment variables
load_dotenv()

# Initialize FastAPI app
app = FastAPI()

# Initialize Supabase client
supabase: Client = create_client(
    supabase_url=os.getenv("SUPABASE_URL"),
    supabase_key=os.getenv("SUPABASE_KEY")
)

# Initialize your OpenAI client
client = OpenAIClient(api_key=os.getenv('OPENAI_API_KEY'),organization=os.getenv('ORG_ID'))

@app.post("/test/batch_processing_result")
async def batch_processing_result(request: Request, background_tasks: BackgroundTasks):
    body = await request.json()
    batch_id = body.get('batch_job_id')
    batch_job = client.batches.retrieve(batch_id)
    # while batch_job.status == 'in_progress':
    batch_job = client.batches.retrieve(batch_id)
    print(batch_job.status)
    # Add the processing task to background tasks
    if batch_job.status == 'completed':
        background_tasks.add_task(process_batch_data, batch_id)
        return {"batch_job_status":'completed'} 
    
    # Immediately return success response
    return {'batch_job_status':'notcompleted'}


async def process_batch_data(batch_id: str):
    try:
        batch_job = client.batches.retrieve(batch_id)
        if batch_job.status == 'completed':
            result_file_id = batch_job.output_file_id
            result = client.files.content(result_file_id).content
            json_str = result.decode('utf-8')
            json_lines = json_str.splitlines()
            
            res = []
            for line in json_lines:
                if line.strip():
                    try:
                        json_dict = json.loads(line)
                        res.append(json_dict)
                    except json.JSONDecodeError as e:
                        print(f"Error decoding JSON on line: {line}\nError: {e}")
            
            for resp in res:
                id = resp.get('custom_id')
                res_id = id.split('-')[1]
                output = json.loads(resp.get('response').get('body').get('choices')[0].get('message').get('content'))
                
                categories = str(output.get('categories'))
                summary = str(output.get('summary'))
                
                supabase_resp = supabase.table("imdb_dataset").select("Description").eq("imdb_id", res_id).execute()
                description = supabase_resp.data[0].get('Description')
                
                insert_response = (
                    supabase.table("imdb_outputs")
                    .insert({
                        "id": res_id, 
                        "description": description,
                        'categories': categories,
                        'summary': summary
                    })
                    .execute()
                )
                print(f"Inserted data for ID: {res_id}")
                
    except Exception as e:
        print(f"Error in background processing: {str(e)}")
        # You might want to log this error or handle it in some way

if __name__ == "__main__":
    uvicorn.run(
        "main:app",
        host="0.0.0.0",
        port=8000,
        reload=True  # Enable auto-reload during development
    )

Utilise the above code to fetch all the data from the batch , that is completed and insert the data to the relevant column into the supabase database.

@app.post("/test/batch_processing_result")
async def batch_processing_result(request: Request, background_tasks: BackgroundTasks):
    body = await request.json()
    batch_id = body.get('batch_job_id')
    batch_job = client.batches.retrieve(batch_id)
    # while batch_job.status == 'in_progress':
    batch_job = client.batches.retrieve(batch_id)
    print(batch_job.status)
    # Add the processing task to background tasks
    if batch_job.status == 'completed':
        background_tasks.add_task(process_batch_data, batch_id)
        return {"batch_job_status":'completed'} 
    
    # Immediately return success response
    return {'batch_job_status':'notcompleted'}

Above snippet actually checks if the batch process is compeleted or not and and then adds the task of inserting the data into the database to a background task , so that api response should not be delayed in the database function and does not face a timeout .

async def process_batch_data(batch_id: str):
    try:
        batch_job = client.batches.retrieve(batch_id)
        if batch_job.status == 'completed':
            result_file_id = batch_job.output_file_id
            result = client.files.content(result_file_id).content
            json_str = result.decode('utf-8')
            json_lines = json_str.splitlines()
            
            res = []
            for line in json_lines:
                if line.strip():
                    try:
                        json_dict = json.loads(line)
                        res.append(json_dict)
                    except json.JSONDecodeError as e:
                        print(f"Error decoding JSON on line: {line}\nError: {e}")
            
            for resp in res:
                id = resp.get('custom_id')
                res_id = id.split('-')[1]
                output = json.loads(resp.get('response').get('body').get('choices')[0].get('message').get('content'))
                
                categories = str(output.get('categories'))
                summary = str(output.get('summary'))
                
                supabase_resp = supabase.table("imdb_dataset").select("Description").eq("imdb_id", res_id).execute()
                description = supabase_resp.data[0].get('Description')
                
                insert_response = (
                    supabase.table("imdb_outputs")
                    .insert({
                        "id": res_id, 
                        "description": description,
                        'categories': categories,
                        'summary': summary
                    })
                    .execute()
                )
                print(f"Inserted data for ID: {res_id}")
                
    except Exception as e:
        print(f"Error in background processing: {str(e)}")
        # You might want to log this error or handle it in some way

Above code snippet fetches the output of the batch , loads that into the json in memory and parses the json and extracts relevant output from the batch output json and inserts the data into the output supabase table.

Each result line also reports its token usage, so you can store the prompt and completion token counts alongside your outputs:

for (ress) in res:
    # print("LLM OUTPUT")
    custom_id = ress.get('custom_id')
    output=ress.get('response').get('body').get('choices')[0].get('message').get('content')
    prompt_tokens = ress.get('response').get('body').get('usage').get('prompt_tokens')
    completion_tokens = ress.get('response').get('body').get('usage').get('completion_tokens')

Conclusion

Batch processing is a reliable way to handle large amounts of data. It lets you automate workflows, streamline data processing and manage costs more effectively.

The process stays the same whatever database you use: set a trigger that fires when you have a significant number of rows, pass that data to an API endpoint that creates a batch job, and save the batch_job_id to a table as a log. Then a cron job checks at regular intervals whether each batch has completed, and a second endpoint collects the results and saves them to your database.

The hussh magazine

Written by Omkar Malpure, and built to read beautifully here — and to travel to 🤫 One on your phone, your glasses, and visionOS, as one immersive magazine you own.

More from the magazine →Back to top ↑

Keep reading

More stories from the magazine

November 27, 2025

OpenAI Public Data Agent Documentation

How hussh’s OpenAI Public Data Agent turns minimal identifiers into enriched JSON profiles for personalization.

July 17, 2025

Voice to Text with Whisper - Let AI Transcribe Anything

Voice is natural. Whether you're dictating notes, talking to a smart speaker, or attending meetings - audio is everywhere. But AI transcription used to be complicated, inaccurate, and expensive.

August 7, 2026

Every Scope Resolves Now: A Systems Review of the Consent Fabric

A full engineering accounting of PCHP and the fabric that serves it: the registry's growth from 47 to 263 scopes, the resolver that went from 8 hand-mapped bindings to resolution by convention, the pseudonym fix that stopped telling subscribers who you are, the economics of a millicent handshake, and a plain ledger of everything that is still not real.

Product

  • Agent One
  • Puppy One
  • Tag One
  • The Store
  • The One Card

Yellow Pages

  • Directory
  • Find an Expert
  • Live Feed
  • Coverage & Markets
  • Connect
  • Ping an Agent

Enterprise

  • For Business
  • RIA Wealth
  • Developers & MCP
  • Partner Portal

Knowledge

  • Guide
  • Podcasts & Transcripts
  • System Architecture

Company

  • About hussh
  • Leadership & Team
  • Media
  • Careers
  • Contact

Trust & Rails

  • PCHP Protocol
  • Data Rights Ledger
  • Zero Telemetry Audit
  • Enclave Hardware
  • Day 0 Trusted Circle
  • The Case for Rights

Private Agent One is free for every American citizen. We do not sell your data, your attention, or your contacts. Certified zero-telemetry by design. Company and product names describe technology interoperability and do not imply endorsement.

Copyright © 2026 Hushh Technologies Corporation. Kirkland, WA. All rights reserved.
Privacy PolicyTerms of UseYour Data RightsAccessibility