Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion .github/workflows/sphinx.yml
Original file line number Diff line number Diff line change
Expand Up @@ -15,7 +15,7 @@ env:
# If these SPHINXOPTS are enabled, then be strict about the
# builds and fail on any warnings.
#SPHINXOPTS: "-W --keep-going -T"
GENERATE_PDF: true # to enable, must be 'true' lowercase
GENERATE_PDF: false # to enable, must be 'true' lowercase
GENERATE_SINGLEHTML: true # to enable, must be 'true' lowercase
PDF_FILENAME: lesson.pdf
MULTIBRANCH: true # to enable, must be 'true' lowercase
Expand Down
78 changes: 78 additions & 0 deletions content/example/download_gfs_mp.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,78 @@
import os
from concurrent.futures import ProcessPoolExecutor, as_completed
import requests
from datetime import date

# Configuration
MAX_PROCESSES = 3 # Adjust based on your connection speed and server limits
DOWNLOAD_DIR = "./downloads" # Directory where to put the files

# Automatically fetch today's date in YYYYMMDD format (e.g., "20260923")
TARGET_DATE = date.today().strftime("%Y%m%d")

# Base URL dynamically injects the current date
BASE_URL = f"https://noaa.gov.{TARGET_DATE}/00/atmos"
FILENAME_TEMPLATE = "gfs.t00z.pgrb2.0p25.f{:03d}"

# Dynamically generate the list of files to download
FILES_TO_DOWNLOAD = [
{
"url": BASE_URL,
"filename": FILENAME_TEMPLATE.format(hour)
}
for hour in range(0, 16, 3)
]


def download_file(file_info):
"""Downloads a single file streaming it in chunks to optimize memory."""
url = file_info["url"]
filename = file_info["filename"]
save_path = os.path.join(DOWNLOAD_DIR, filename)

try:
# Stream the download to avoid loading huge files into memory all at once
with requests.get(os.path.join(url,filename), stream=True, timeout=15) as response:
response.raise_for_status() # Check for HTTP errors (404, 500, etc.)

with open(save_path, "wb") as file:
for chunk in response.iter_content(
chunk_size=8192
): # 8KB chunks
if chunk:
file.write(chunk)

return f"Successfully downloaded: {filename}"

except requests.exceptions.RequestException as e:
return f"Failed to download {filename}. Error: {e}"
except Exception as e:
return f"An unexpected error occurred for {filename}: {e}"


def main():
# Ensure destination directory exists
os.makedirs(DOWNLOAD_DIR, exist_ok=True)

print(
f"Starting download of {len(FILES_TO_DOWNLOAD)} files using {MAX_PROCESSES} processes...\n"
)

# Use ProcessPoolExecutor to handle concurrent downloads
with ProcessPoolExecutor(max_workers=MAX_PROCESSES) as executor:
# Submit all tasks to the executor
future_to_url = {
executor.submit(download_file, file): file
for file in FILES_TO_DOWNLOAD
}

# Process results as they finish
for future in as_completed(future_to_url):
result_message = future.result()
print(result_message)

print("\nAll download processes completed.")


if __name__ == "__main__":
main()
79 changes: 79 additions & 0 deletions content/example/download_gfs_mt.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,79 @@
import os
from concurrent.futures import ThreadPoolExecutor, as_completed
import requests
from datetime import date

# Configuration
MAX_THREADS = 3 # Adjust based on your connection speed and server limits
DOWNLOAD_DIR = "./downloads" # Directory where to put the files

# Automatically fetch today's date in YYYYMMDD format (e.g., "20260923")
TARGET_DATE = date.today().strftime("%Y%m%d")

# Base URL dynamically injects the current date
BASE_URL = f"https://noaa.gov.{TARGET_DATE}/00/atmos"
FILENAME_TEMPLATE = "gfs.t00z.pgrb2.0p25.f{:03d}"

# Dynamically generate the list of files to download
FILES_TO_DOWNLOAD = [
{
"url": BASE_URL,
"filename": FILENAME_TEMPLATE.format(hour)
}
for hour in range(0, 16, 3)
]



def download_file(file_info):
"""Downloads a single file streaming it in chunks to optimize memory."""
url = file_info["url"]
filename = file_info["filename"]
save_path = os.path.join(DOWNLOAD_DIR, filename)

try:
# Stream the download to avoid loading huge files into memory all at once
with requests.get(os.path.join(url,filename), stream=True, timeout=15) as response:
response.raise_for_status() # Check for HTTP errors (404, 500, etc.)

with open(save_path, "wb") as file:
for chunk in response.iter_content(
chunk_size=8192
): # 8KB chunks
if chunk:
file.write(chunk)

return f"Successfully downloaded: {filename}"

except requests.exceptions.RequestException as e:
return f"Failed to download {filename}. Error: {e}"
except Exception as e:
return f"An unexpected error occurred for {filename}: {e}"


def main():
# Ensure destination directory exists
os.makedirs(DOWNLOAD_DIR, exist_ok=True)

print(
f"Starting download of {len(FILES_TO_DOWNLOAD)} files using {MAX_THREADS} threads...\n"
)

# Use ThreadPoolExecutor to handle concurrent downloads
with ThreadPoolExecutor(max_workers=MAX_THREADS) as executor:
# Submit all tasks to the executor
future_to_url = {
executor.submit(download_file, file): file
for file in FILES_TO_DOWNLOAD
}

# Process results as they finish
for future in as_completed(future_to_url):
result_message = future.result()
print(result_message)

print("\nAll download processes completed.")


if __name__ == "__main__":
main()
88 changes: 88 additions & 0 deletions content/example/duckdb_mt.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,88 @@
import duckdb
from concurrent.futures import ThreadPoolExecutor, as_completed
import threading
import random

# ==========================================
# 1. Setup Database Connection & Table
# ==========================================
duckdb_con = duckdb.connect('my_persistent_db.duckdb')

duckdb_con.execute("""
CREATE OR REPLACE TABLE my_inserts (
thread_name VARCHAR,
insert_time TIMESTAMP DEFAULT current_timestamp
)
""")


# ==========================================
# 2. Worker Tasks (Writer & Reader)
# ==========================================
def write_task(duckdb_con, task_id):
# Create a unique, thread-safe cursor from the main connection
local_con = duckdb_con.cursor()

# Track the underlying OS thread name alongside the task ID
thread_name = f"writer_task_{task_id} ({threading.current_thread().name})"

local_con.execute("""
INSERT INTO my_inserts (thread_name) VALUES (?)
""", (thread_name,))
return f"Task {task_id} successfully inserted data."


def read_task(duckdb_con, task_id):
# Create a unique, thread-safe cursor from the main connection
local_con = duckdb_con.cursor()

thread_name = f"reader_task_{task_id} ({threading.current_thread().name})"
results = local_con.execute("""
SELECT ? AS worker, count(*) AS row_counter, current_timestamp
FROM my_inserts
""", (thread_name,)).fetchall()

return results


# ==========================================
# 3. Queue Tasks and Execute with ThreadPool
# ==========================================
write_task_count = 500
read_task_count = 25000

# Prepare the collection of callable functions and their arguments
tasks = []

for i in range(write_task_count):
tasks.append((write_task, (duckdb_con, i)))

for j in range(read_task_count):
tasks.append((read_task, (duckdb_con, j)))

# Shuffle tasks to mix reader and writer execution order randomly
random.seed(6)
random.shuffle(tasks)

# Execute via ThreadPoolExecutor
# max_workers can be tuned based on your machine's hardware capabilities
with ThreadPoolExecutor(max_workers=5) as executor:
# Submit all shuffled operations to the thread pool
futures = [executor.submit(func, *args) for func, args in tasks]

# Process results dynamically as they complete
for future in as_completed(futures):
try:
result = future.result()
## If it's a reader task (returns a list), print it out
#if isinstance(result, list):
# print(f"Reader Result: {result}")
except Exception as e:
print(f"A task generated an exception: {e}")


# ==========================================
# 4. Show Final Results
# ==========================================
print("\n--- Final Database Contents ---")
print(duckdb_con.execute("SELECT * FROM my_inserts ORDER BY insert_time").df())
25 changes: 25 additions & 0 deletions content/example/integration_serial.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
import math
import time

# Grid size
n = 100000000

def integration_serial(n):
h = 1.0 / float(n)
mysum = 0.0

for i in range(n):
x = h * (i + 0.5)
mysum += x ** (3/2)

return h * mysum

if __name__ == "__main__":
starttime = time.time()
integral = integration_serial(n)
endtime = time.time()

print("Integral value is %e, Error is %e" % (integral, abs(integral - 2/5))) # The correct integral value is 2/5
print("Time spent: %.2f sec" % (endtime-starttime))

# 13.63 sec
41 changes: 41 additions & 0 deletions content/example/kommun.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,41 @@
import numpy as np
import pandas as pd
import geopandas as gpd
from shapely.geometry import Point
import time
from concurrent.futures import ThreadPoolExecutor

# load data
points = pd.read_csv("./municipality.csv",usecols=["Locality", "Latitude", "Longitude"])
polygons = gpd.read_file("./kommun_se.geojson")

n_points = len(points)
n_polygons = len(polygons)

def check_polygon(polygon_idx):
"""Worker function to process a single polygon"""
current_polygon = polygons.iloc[polygon_idx]["geometry"]
out_points = []
# manually loop over all points, check if polygon contains that point
for i in range(n_points):
current_point = points.iloc[i, :]
# Note: Shapely Point expects (Longitude, Latitude) i.e., (x, y)
if current_polygon.contains(Point(current_point.Longitude,current_point.Latitude)):
out_points.append(current_point.Locality)
return out_points



if __name__ == "__main__":
print(f"Starting with single thread...")

t_start=time.time()

points_per_polygon = {}

for polygon_idx in range(n_polygons):
points_per_polygon[polygon_idx] = check_polygon(polygon_idx)

t_end=time.time()

print(f"Finish Processing {n_points} points against {n_polygons} polygons in {t_end - t_start:.2f} seconds.")
29 changes: 29 additions & 0 deletions content/example/omp_test.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
import numpy as np
import time

def timming_matrix_inversion(size: int = 4000) -> float:
"""Generates a symmetric matrix and measures the time taken to invert it."""
# Create a random matrix
A = np.random.random((size, size))

# Make it symmetric (Note: A * A.T is element-wise multiplication.
# If you meant matrix multiplication, use A @ A.T or np.dot(A, A.T))
A = A * A.T

# Measure inversion time
time_start = time.time()
np.linalg.inv(A)
time_end = time.time()

return time_end - time_start


def main() -> None:
matrix_size = 4000
duration = timming_matrix_inversion(matrix_size)

print(f"Time spent for inverting a random matrix with a size of ({matrix_size}x{matrix_size}) is {round(duration, 2)} s")


if __name__ == "__main__":
main()
23 changes: 23 additions & 0 deletions content/example/race.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
from multiprocessing import Value

# define a function to increment the value by 1
def inc(i):
val.value += 1

# using a large number to see the problem
n = 100000

# create a shared data and initialize it to 0
val = Value('i', 0)
with ThreadPoolExecutor(max_workers=4) as pool:
pool.map(inc, range(n))

print(val.value)

# create a shared data and initialize it to 0
val = Value('i', 0)
with ProcessPoolExecutor(max_workers=4) as pool:
pool.map(inc, range(n))

print(val.value)
Loading
Loading