How to async call API one million times with different inputs using Python

Viewed 376

I have a list of IDs (about 1 Million) which I want to iterate over and make async calls to an API endpoint and store the response in a file.

So far I have searched many ways but they all make calls only using a single ID, couldn't find a way to call the same APIs URL effectively with different IDs.

I understood the performance impact so I want to make use a combination of asyncio, aiohttp, utilizing CPU cores and threads efficiently to make parallel calls. (concurrent, multiprocessing)

My Plan is to split the big list into small chunks of 10 IDs are send them to different threads.

Sample data: id_data = ['0', '1', '2', '3', '4', '5', '6', '7', '8', '9', '10', '11', '12', '13', '14', '15', '16', '17', '18', '19', '20', '21', '22', '23', '24', '25', '26', '27', '28', '29'] and so-on till millions...

API Endpoint = .get('https://random.url.of.web/v70/services/each Id comes here')

My implementation as of now:

id_data = ['0', '1', '2', '3', '4', '5', '6', '7', '8', '9', '10', '11', '12', '13', '14', '15', '16', '17', '18', '19', '20', '21', '22', '23', '24', '25', '26', '27', '28', '29']

def divide_chunks(l, n): 
    for i in range(0, len(l), n):  
        yield l[i:i + n]

async def get_and_scrape_pages(num_pages: list,output_file: str):
    async with \
    aiohttp.ClientSession(headers=headers) as client, \
    aiofiles.open(output_file, "a+", encoding="utf-8") as f:
        for z in range(len(num_pages)):
            async with client.get('https://random.url.of.web/v70/services/'+id_data[z]) as response:
                print(z+1,accid_data[z],response.status)
                page = await response.text() + "-----" + str(z+1)
                #print(page)
                
                await f.write(page + "\n")
        await f.write("\n")

def start_scraping(num_pages: list,output_file: str):
    print("\n\tscraping now...\n")
    x = list(divide_chunks(accid_data, 10)) 
    for p in range(len(x)):
        asyncio.run(get_and_scrape_pages(x[p],output_file))


def main():
    NUM_API = 30
    NUM_CORES = cpu_count() # Our number of CPU cores (including logical cores)
    OUTPUT_FILE = "C://Users//40102046//eclipse-workspace//api_extraction//logs_1.txt" # File to append our scraped titles to
    
    print("number of CPU cores (including logical cores)",NUM_CORES)
    PAGES_PER_CORE = floor(NUM_API / NUM_CORES)
    PAGES_FOR_FINAL_CORE = PAGES_PER_CORE + NUM_API % PAGES_PER_CORE # For our final core
    
    print("PAGES_PER_CORE",PAGES_PER_CORE)
    print("PAGES_FOR_FINAL_CORE",PAGES_FOR_FINAL_CORE)
    
    futures = []
    
    with concurrent.futures.ProcessPoolExecutor(NUM_CORES) as executor:
        for i in range(NUM_CORES): 
            new_future = executor.submit(
                start_scraping(),
                num_pages=PAGES_PER_CORE,
                output_file=OUTPUT_FILE,
            )
            print("from main def:",i+1)
            futures.append(new_future)

        futures.append(
            executor.submit(
                start_scraping, PAGES_FOR_FINAL_CORE, OUTPUT_FILE
            )
        )

    concurrent.futures.wait(futures)
  
if __name__ == "__main__":
    main()

This code: always sends the first list of 10 "IDs" to all the threads every time. I want to send each nested list separately to different threads, any help or guidance is appreciated.

1 Answers

Here are some of the mistakes I think are causing the issue you're dealing with, along with some recommendations.

  1. z ought to iterate over the values of num_pages, not the range of its length. And this should be renamed to something like id_chunk to reflect its contents.
       for z in chunk_ids:
            async with client.get(
                    'https://random.url.of.web/v70/services/' + z) as response:
                print(z, response.status)
                page = await response.text() + "-----" + str(z)
                # print(page)
                await f.write(page + "\n")
  1. I would move the computation of n into the divide_chunks function, making num_chunks (i.e. the number of cores) a parameter to this function. You'll want to ensure the divide_chunks function returns all of the chunks, including the very last one that may be smaller than n.
  2. I would double check your code in main, as I currently don't see the work being divided among the calls to start_scraping. Within main, you need to specify the chunk on which the call to start_scraping ought to operate, and invoke divide_chunks in main instead of start_scraping. None of the calls from start_scraping should touch global variables. The chunk they're concerned with is locally available, and the call to get_and_scrape_pages is provided the data it needs as a parameter.
  3. The implementation looks more complicated than it needs to be. For instance, I don't think you're going to gain much in performance by running the file I/O asynchronously. I would also discourage the use of global variables in general.
Related