Programming languages allow us to communicate with computers, and they operate like sets of instructions. There are numerous types of languages, including procedural, functional, object-oriented, and more. Whether you’re looking to learn a new language or trying to find some tips or tricks, the resources in the Languages Zone will give you all the information you need and more.
Embabel vs LangGraph4j: Two Agentic Philosophies for Investment and Risk Analysis in BFSI
How Go Maps Work: From Buckets to Swiss Tables
When many developers think about recommendation engines, they think of machine learning: collaborative filtering models, matrix factorization, embedding vectors, and training pipelines. What surprises many people is that you can build a genuinely useful recommendation system with nothing more than a graph database and several Cypher queries. No scikit-learn, no TensorFlow, no model training. Just the natural structure of the data doing the work. In this article, we'll build a product recommendation engine on top of Neo4j Aura using two Jupyter notebooks. The first generates a realistic synthetic dataset and loads it into Aura. The second runs four recommendation queries directly in Cypher and visualizes the results with Plotly. Everything runs locally in a Python virtual environment against a free cloud Neo4j instance. The full source code is available on GitHub. Why Graphs Are a Natural Fit for Recommendations The core intuition behind most recommendation approaches is relationship: this customer bought that product, those products appear together in the same order, this product shares attributes with that one. In a relational database, capturing these relationships means multiple self-joins across large tables. A query like "find products bought by customers who also bought what this customer bought" quickly becomes difficult to write and expensive to execute at scale. In a graph, that same question is a traversal. We follow edges from a customer to the products they purchased, hop across to other customers who share those products, and collect what else those customers bought. The query is short, the intent is clear, and the graph engine is optimized for exactly this kind of path-following work. Prerequisites AuraDB is Neo4j's fully managed cloud database. A free tier is available with no credit card required. Sign up at Get Started for Free.Create a new AuraDB Free instance.When the instance is created, download or note the credentials — the connection URI, username, and password.Once the instance is running, open the Query tab and connect to the instance.Confirm it's empty with MATCH (n) RETURN count(n) which should return 0 A virtual environment is highly recommended. For example: Shell python3 -m venv ~/recommendation-engine-env source ~/recommendation-engine-env/bin/activate Before starting Jupyter, export the connection details as environment variables in your shell: Shell export NEO4J_URI="neo4j+s://xxxx.databases.neo4j.io" export NEO4J_USERNAME="your_username_here" export NEO4J_PASSWORD="your_password_here" The Graph Model Before we write any code, let's define the graph. We have four node types and three relationship types. Nodes Customer – id, name, email, city, country.Product – id, name, description, price.Category – name (e.g., Electronics, Clothing, Books).Tag – name (e.g. "wireless", "eco-friendly", "premium"). Relationships (:Customer)-[:PURCHASED {order_id, quantity, order_date}]->(:Product) — order metadata lives on the relationship rather than a separate Order node, which keeps our Cypher clean.(:Product)-[:BELONGS_TO]->(:Category)(:Product)-[:TAGGED_WITH]->(:Tag) The decision to put order_id, quantity and order_date on the PURCHASED relationship is worth discussing. It means a single customer can have multiple PURCHASED relationships to the same product (each with a different order_id) and we can group by order_id to find products that appeared together in the same basket — which is exactly what our co-purchase query needs. Figure 1 illustrates exactly this point, as we have a customer, two products, and the same order_id. Figure 1. Shared order_id enables co-purchase queries Notebook 1: Data Generation and Loading Rather than sourcing an external dataset, we'll generate synthetic data using Faker. This keeps the notebook fully self-contained, and readers can run it as-is without downloading anything. We'll generate 2,000 customers, 500 products across 15 categories, and 20,000 orders. Each order is a basket of several products sharing the same order_id — this is the key design decision that makes the frequently-bought-together query work. With an average basket of 3 products, we end up with around 60,000 PURCHASED relationships in the graph. Realistic Product Names Faker's default catch_phrase() method produces output like "Proactive exuding encoding" — readable enough for a demo but not really useful in an article. Instead, we define a PRODUCT_VOCAB dictionary keyed by category, each containing lists of adjectives, nouns, use cases, and benefit statements. A product name is then a simple combination, as follows: Python def make_product_name(category): vocab = PRODUCT_VOCAB[category] adj = random.choice(vocab["adjectives"]) noun = random.choice(vocab["nouns"]) return f"{adj} {noun}" def make_product_description(category, name): vocab = PRODUCT_VOCAB[category] use_case = random.choice(vocab["use_cases"]) benefit = random.choice(vocab["benefits"]) return f"The {name} is designed for {use_case}. {benefit}." This gives us names like "Wireless Noise-Canceling Earbuds," "Organic Ground Coffee" and "Ergonomic Lumbar Support Cushion" — realistic enough to make the recommendation output meaningful. Basket-Based Order Generation Each order picks a random customer, generates a unique order_id, and samples several products into a basket. We then flatten the basket into individual order lines, each carrying the shared order_id: Python orders = [] for _ in range(NUM_ORDERS): order_id = str(uuid.uuid4()) customer = random.choice(customers) order_date = (start_date + timedelta(days=random.randint(0, 730))).strftime("%Y-%m-%d") basket = random.sample(products, k=random.randint(2, 4)) for product in basket: orders.append({ "order_id": order_id, "customer_id": customer["id"], "product_id": product["id"], "quantity": random.randint(1, 5), "order_date": order_date }) Loading Into Aura Data loading is in batches of 100 using MERGE statements. To show progress during the load, we'll use tqdm as ~60,000 order lines can take several minutes, and the progress bars make it easy to see what's happening: Python with driver.session() as session: customer_batches = range(0, len(customers), BATCH_SIZE) for i in tqdm(customer_batches, desc="Loading customers", unit="batch", colour="#1f77b4"): session.execute_write(load_customers, customers[i:i+BATCH_SIZE]) product_batches = range(0, len(products), BATCH_SIZE) for i in tqdm(product_batches, desc="Loading products ", unit="batch", colour="#1f77b4"): session.execute_write(load_products, products[i:i+BATCH_SIZE]) for product_id, tags in tqdm(product_tags.items(), desc="Loading tags ", unit="product", colour="#1f77b4"): session.execute_write(load_tags, product_id, tags) order_batches = range(0, len(orders), BATCH_SIZE) for i in tqdm(order_batches, desc="Loading orders ", unit="batch", colour="#1f77b4"): session.execute_write(load_orders, orders[i:i+BATCH_SIZE]) A verification query at the end confirms the counts. The Four Recommendation Queries Notebook 2 runs four Cypher queries against the loaded graph, each implementing a different recommendation strategy. Before running any query, we fetch a stable seed customer, product, and category: Python with driver.session() as session: customer = session.run(""" MATCH (c:Customer) RETURN c.id AS customer_id, c.name AS customer_name ORDER BY c.name ASC LIMIT 1 """).single() product = session.run(""" MATCH (p:Product)<-[r:PURCHASED]-() RETURN p.id AS product_id, p.name AS product_name, count(r) AS order_count ORDER BY order_count DESC LIMIT 1 """).single() top_cat = session.run(""" MATCH (p:Product)-[:BELONGS_TO]->(cat:Category) RETURN cat.name AS category, count(p) AS total ORDER BY total DESC LIMIT 1 """).single() We pick the alphabetically first customer for consistency, the most-purchased product to ensure co-purchase data exists, and the category with the most products for the trending query. This makes the notebook reproducible across runs. Query 1: Collaborative Filtering The classic "customers who bought this also bought" approach. We find customers who share at least one purchased product with the seed customer, then collect what else those customers bought — excluding anything the seed customer already purchased. Python def collaborative_filtering(tx, customer_id, limit=5): result = tx.run(""" MATCH (target:Customer {id: $customer_id})-[:PURCHASED]->(p:Product) <-[:PURCHASED]-(other:Customer)-[:PURCHASED]->(rec:Product) WHERE NOT (target)-[:PURCHASED]->(rec) RETURN rec.id AS id, rec.name AS product, rec.price AS price, count(other) AS score ORDER BY score DESC, id ASC LIMIT $limit """, customer_id=customer_id, limit=limit) return result.data() The score is the number of other customers whose purchasing overlap with our target customer also led them to buy the recommended product. A higher score means more customers in the overlap group bought it, making it a stronger signal. In Cypher, the traversal reads almost like the description: start at the target customer, follow PURCHASED edges to products, hop to other customers who bought the same products, then follow their PURCHASED edges to new products. Query 2: Frequently Bought Together This query finds products that appeared in the same order as the seed product. The key is matching on order_id across two PURCHASED relationships from the same customer: Python def frequently_bought_together(tx, product_id, limit=5): result = tx.run(""" MATCH (p:Product {id: $product_id})<-[r1:PURCHASED]-(c:Customer) -[r2:PURCHASED]->(other:Product) WHERE r1.order_id = r2.order_id AND other.id <> $product_id RETURN other.id AS id, other.name AS product, other.price AS price, count(c) AS frequency ORDER BY frequency DESC, id ASC LIMIT $limit """, product_id=product_id, limit=limit) return result.data() The WHERE r1.order_id = r2.order_id clause is what makes this work. It constrains the traversal to only consider cases where both products were part of the same order, not just bought by the same customer at different times. frequency counts how many distinct customers placed an order containing both products together. Query 3: Content-Based Filtering Rather than looking at purchase behavior, this query finds products similar to the seed product based on shared tags. The more tags two products have in common, the more similar they are: Python def content_based(tx, product_id, limit=5): result = tx.run(""" MATCH (p:Product {id: $product_id})-[:TAGGED_WITH]->(t:Tag) <-[:TAGGED_WITH]-(rec:Product) WHERE rec.id <> $product_id RETURN rec.id AS id, rec.name AS product, rec.price AS price, count(t) AS shared_tags ORDER BY shared_tags DESC, id ASC LIMIT $limit """, product_id=product_id, limit=limit) return result.data() The traversal goes outward from the seed product through its tags, then back inward to any other product that shares those same tags. count(t) gives the number of shared tags, which serves as a simple but effective similarity score. This approach works without any purchase history, making it useful for recommending products to new customers or for newly listed products with no order data yet. Query 4: Trending in Category This query finds the most purchased products in the top category within a fixed date window. In our case, this is from 2024-10-01 onwards: Python def trending_in_category(tx, category_name, cutoff="2024-10-01", limit=5): result = tx.run(""" MATCH (p:Product)-[:BELONGS_TO]->(cat:Category {name: $category_name}) MATCH (:Customer)-[r:PURCHASED]->(p) WHERE date(r.order_date) >= date($cutoff) RETURN p.id AS id, p.name AS product, p.price AS price, count(r) AS purchases ORDER BY purchases DESC, id ASC LIMIT $limit """, category_name=category_name, cutoff=cutoff, limit=limit) return result.data() We use date() conversion on the stored string order_date to enable date comparison. count(r) counts individual PURCHASED relationships rather than distinct customers, so a customer who bought the same product multiple times within the window is counted each time — reflecting genuine demand volume rather than unique buyer count. Notebook 2: Results Each query outputs a table followed by a Plotly horizontal bar chart. Here are the results for our seed data. Collaborative Filtering Figure 2 returns five products. The top recommendation is An Introduction to Public Speaking, driven by the number of customers whose purchasing overlap with Aaron Boyd also led them to buy it. Heavy-Duty Cable Management Box and Educational Coding Robot follow closely, showing that the overlap group bought broadly across categories rather than clustering in one area. Figure 2. Collaborative filtering Frequently Bought Together Figure 3 shows products co-purchased with the Durable Grooming Brush in the same order basket. The top results — Waterproof Hammock and Natural Body Lotion at frequency 4, followed by Adjustable Lumbar Support Cushion, Sugar-Free Collagen Powder and Slim-Fit Hiking Vest at frequency 3 — show which products most commonly appeared alongside the seed product in the same order. The cross-category spread here (Beauty, Outdoor, Clothing, Health, Office) is a feature of random synthetic data; in a real system, we'd expect more category clustering. Figure 3. Frequently bought together Content-Based Filtering Figure 4 finds products sharing the most tags with the seed product. All five results share 2 tags with the Durable Grooming Brush — Smart Mechanical Keyboard, Waterproof Toiletry Bag, Ergonomic Whiteboard, Cold-Pressed Hot Sauce, and Durable Dumbbell Pair. The cross-category reach (Sports, Food & Drink, Office, Travel, Electronics) illustrates the tag graph doing its job: shared attributes like "durable" or "waterproof" create similarity links that cross category boundaries, which is useful for surface-level discovery recommendations. Figure 4. Content-based filtering Trending in Category Figure 5 shows the top 5 products in Toys — the category with the most products in our graph — with purchase counts from 2024-10-01 onwards. Battery-Free Coding Robot leads, followed by Battery-Free Building Blocks Set, Interactive Remote Control Car, Wooden Magnetic Drawing Board, and Creative Puzzle Game. The scores are tight here, which makes sense because within a single category over a fixed time window, popular products tend to cluster around similar purchase volumes. Figure 5. Trending Summary We've built a working product recommendation engine using nothing but Neo4j, Cypher, and a few Python libraries. No ML framework, no training data, no model deployment. The four queries cover the most common recommendation patterns in production systems: Collaborative filteringCo-purchase analysisContent similarityTrending detection The graph model is the foundation that makes this possible. Storing orders as relationships with properties means co-purchase queries are a natural traversal rather than a complex join. Adding tags as nodes means similarity queries are just path-matching. Because everything lives in the same graph, we can also combine these approaches. For example, filtering collaborative filtering results by tag similarity using a single extended Cypher query. The full source code is available on GitHub.
Unit testing business logic often requires surprisingly little business data. Suppose we want to test a program that loads an order, calculates its total, and rejects it when the amount exceeds a limit. The decision we want to verify is simple: The program loads the specified order.It asks for the order’s total.If the total is too high, it does not approve the order.It returns FALSE. Yet a conventional Java unit test may have to construct an Order. That order may require a customer, line items, currencies, prices, tax information, identifiers, and other objects that have nothing to do with the decision being tested. Builders, fixtures and mocking frameworks reduce the typing, but they do not eliminate the underlying problem: the test must participate in the internal representation of the domain model. BUBAS takes a different approach. Domain objects are opaque to a BUBAS program. Because the program cannot inspect them, a unit test does not need to construct them. It needs only a token. The Business Program Consider this BUBAS program: SQL PROGRAM ApproveOrder(orderId INTEGER, limit DECIMAL) RETURNS BOOLEAN DECLARE purchase Order DECLARE total DECIMAL purchase = LOAD_ORDER(orderId) IF NOT ORDER_WAS_FOUND(purchase) THEN LOG_EVENT "ERROR", "no such order: " + orderId RETURN FALSE END IF total = ORDER_TOTAL(purchase) IF total > limit THEN LOG_EVENT "INFO", "over limit: " + total RETURN FALSE END IF APPROVE purchase RETURN TRUE END. Order is a Java domain type registered by the application embedding BUBAS. The program can store an Order in a variable and pass it to operations that accept an Order, but it cannot access its fields or invoke its methods. There is no expression such as: SQL purchase.customer.account.balance If the program needs information about an order, the application must expose an operation for obtaining it: SQL total = ORDER_TOTAL(purchase) This restriction is primarily an encapsulation mechanism. The business program depends on the vocabulary of its domain rather than on the internal structure of Java objects. It also has an important consequence for testing. Replace the Object With Identity Here is a BUNIT test for the over-limit case: Gherkin PROGRAM OverLimitIsRejected "LOAD_ORDER" WITH ARGS(42) RETURNS "o1" "ORDER_TOTAL" WITH ARGS("o1") RETURNS 1500.00 "APPROVE _" IS MOCKED ARGUMENT "orderId" IS 42 ARGUMENT "limit" IS 1000.00 RUN RESULT IS FALSE "APPROVE _" WAS NOT CALLED END. The string "o1" is not an order serialized as text. It does not contain an order number, a total or any other property. It is a test token representing one opaque Order. The first mock says: Plain Text "LOAD_ORDER" WITH ARGS(42) RETURNS "o1" When the program calls LOAD_ORDER(42), BUNIT returns the token "o1" in place of the real Java object. The program stores it in purchase. Later it calls: Java ORDER_TOTAL(purchase) The second mock recognizes that same token and returns 1500.00. The program cannot tell that "o1" is not a real Order. It has no operation with which to inspect the object. It can only pass the value back through the vocabulary supplied by the host application. For this test, identity is all the domain object needs. We Are Testing the Conversation A BUBAS business program contains decisions and orchestration. Algorithms, persistence, infrastructure, and domain-object implementations remain in Java. Its unit test should therefore concentrate on questions such as: Which domain operations were invoked?With what arguments?What values did those operations return?Which branch did the program select?Which operations were deliberately not invoked?What result did the program produce? In the example, we do not test how ORDER_TOTAL calculates a total. That belongs in the Java test for the implementation of ORDER_TOTAL. We test what the business program does when ORDER_TOTAL reports 1500.00. This division gives us two focused tests rather than one oversized test: Java tests verify the individual domain operations.BUNIT tests verify how a business program coordinates them. The BUNIT test documents the business scenario directly. An order identified by 42 exists, its total is 1500.00, the approval limit is 1000.00, and the program must not approve it. The test does not explain how to manufacture an object graph that produces those facts. Opacity Buys Mockability Mocking domain objects in a general-purpose language is often difficult precisely because the production code can observe so much about them. It may call methods, inspect nested objects, compare values, serialize the object or pass it to code that expects a particular implementation. A substitute must reproduce every observable property used along the tested path. An opaque BUBAS value has only the observations provided by the registered vocabulary. If the vocabulary exposes ORDER_TOTAL, then the mock controls the answer to ORDER_TOTAL. If it does not expose the customer’s internal account object, neither the program nor the test needs to know that such an object exists. The object boundary and the testing boundary are the same boundary. This is stronger than merely saying that business programs should avoid inspecting domain objects. They cannot inspect them unless the embedder deliberately provides an operation that does so. Consequently, a token can stand in for any opaque value as long as the mocks define how the exposed operations respond to it. Multiple objects require only multiple identities: Plain Text "LOAD_ORDER" WITH ARGS(42) RETURNS "o1" "LOAD_ORDER" WITH ARGS(43) RETURNS "o2" "ORDER_TOTAL" WITH ARGS("o1") RETURNS 1500.00 "ORDER_TOTAL" WITH ARGS("o2") RETURNS 200.00 The test describes the distinctions that matter without constructing either order. The Test Uses the Real Language A dangerous form of mocking creates a second, simplified interface used only by tests. Eventually, the production vocabulary changes while the test vocabulary does not. BUNIT does not compile the business program against a parallel language. The program under test is compiled against the real sealed BUBAS language. Mocking happens later, at dispatch. Therefore, the test cannot silently keep using an operation that no longer exists in the production language. Nor can it casually return a value of the wrong BUBAS type. Before executing a test, BUNIT checks the mocks and the test configuration. It can report problems such as: A mock declared with the wrong number of arguments;A mock returning a value incompatible with the real operation;An argument supplied for a parameter the program does not accept;A mocked command that should initialize a variable but does not provide its value. The test reports these errors before the business program runs. The test remains artificial — as every unit test is — but it is artificial inside the actual language contract. Do Not Assert Everything A test becomes fragile when it records every interaction, whether or not that interaction matters to the scenario. BUNIT allows an expectation to specify only the relevant part of a call. For example: Plain Text "LOG_EVENT _, _" WAS CALLED WITH ARGS("INFO", CONTAINS("over limit")) The test requires an informational log message containing "over limit". It does not require the complete message to remain byte-for-byte identical. Similarly: Plain Text "APPROVE _" WAS NOT CALLED expresses the important negative requirement without inventing an Order merely to compare it with another Order. The purpose is not to reproduce the execution trace. It is to state the observable facts that define the business case. What This Does Not Test Opaque tokens do not prove that the Java implementation of LOAD_ORDER returns the right order. They do not prove that ORDER_TOTAL calculates taxes correctly or that APPROVE commits a transaction. Those operations require their own Java unit and integration tests. BUNIT tests the program at the orchestration boundary. This makes it possible to test business decisions without databases, service containers, or complete domain-object graphs, but it does not replace testing below or beyond that boundary. Nor does BUNIT make every Java application automatically testable. The application developer first has to expose a suitably designed vocabulary. If one enormous operation performs loading, calculation, approval and notification internally, BUNIT can mock that operation but cannot test the decisions hidden inside it. Testability therefore provides feedback about vocabulary design. Operations should represent meaningful domain capabilities at the level where business programs genuinely make choices. The Deeper Result Opaque domain types may initially look like a limitation. The program cannot examine its own values freely. It has to ask the vocabulary to interpret them. That limitation creates a clean separation: Java owns domain representation and implementation.BUBAS owns orchestration and decisions.BUNIT replaces domain capabilities at that same boundary.Tokens replace complex objects with identity when identity is all the test requires. The production program becomes independent of domain-object structure. The unit test inherits that independence. We do not need a fake Order with a fake customer containing fake line items whose prices happen to add up to 1500.00. For this business decision, we need only to say: Plain Text "ORDER_TOTAL" WITH ARGS("o1") RETURNS 1500.00 The business program never needed to know what was inside the order. Neither does its test. The detailed code and the BUBAS framework are available as open source at https://github.com/verhas/bubas.
One of the most important parts of API test automation is validating the response body to ensure data integrity. This step plays a key role in functional API testing, as it helps confirm that the API is returning the right data in the expected format. Response body validation isn’t limited to a specific request type; it applies equally to POST, GET, PUT, and PATCH APIs. The same validation approach can be used for any API response to verify the data returned by the service. Playwright offers multiple ways to validate response bodies. In this tutorial, I’ll walk you through these approaches to help you efficiently perform assertions on the response data using best practices. Checkout the previous tutorial blog to learn about Installation, the demo application, and how to send GET API requests with Playwright. How to Verify the Response Structure Response structure checks ensure that an API consistently returns data in the expected format, protecting the contract between backend services and their consumers. They help catch breaking changes early, such as missing or renamed fields, even when the API still returns a successful status code. TypeScript test("GET Order details and perform structure check", async ({ request }) => { const response = await request.get("http://localhost:3004/getOrder/", { params: { user_id: "1", }, failOnStatusCode: true, }); const responseBody = await response.json(); expect(responseBody).toHaveProperty("message"); expect(responseBody).toHaveProperty("orders"); expect(responseBody.orders[0]).toHaveProperty("id"); expect(responseBody.orders[0]).toHaveProperty("product_name"); }); This test focuses on validating the structure of the API response. It validates that the response body contains the expected top-level keys and that each order object includes the required fields. Basic Assertions The basic assertions validate API success and data presence, making them a good first layer of verification before deeper structure or data-level checks. TypeScript test("Get order details and perform basic level verification", async ({ request, }) => { const response = await request.get("http://localhost:3004/getOrder/", { params: { user_id: 1, }, failOnStatusCode: true, }); const responseBody = await response.json(); expect(responseBody.message).toBe("Order found!!"); expect(Array.isArray(responseBody.orders)).toBeTruthy(); expect(responseBody.orders.length).toBeGreaterThan(0); }); This test performs a basic level check to confirm that the endpoint works as expected and returns the expected data in the response. After parsing the response body, the assertions focus on the following essential basic-level checks: TypeScript expect(responseBody.message).toBe("Order found!!"); The above line of code verifies that the API returns the expected message text in the response body. TypeScript expect(Array.isArray(responseBody.orders)).toBeTruthy(); This line of code ensures that the orders field in the response is an array, validating the basic response format. TypeScript expect(responseBody.orders.length).toBeGreaterThan(0); This part of the test confirms that at least one order is returned in the orders array, ensuring the response contains required data. How to Verify Response Data With Details Validating the actual data returned in the response is essential to ensure that the API response contains the correct values. TypeScript test("Get order and verify order details", async ({ request }) => { const response = await request.get("http://localhost:3004/getOrder/", { params: { user_id: "1", }, failOnStatusCode: true, }); const responseBody = await response.json(); const order = responseBody.orders[0]; expect(order.id).not.toBeNull(); expect(order.id).toBeDefined(); expect(order.user_id).toEqual("1"); expect(order.product_id).toEqual("79"); expect(order.product_name).toEqual("5 star 10gm Chocobar"); }); The following code ensures that the response has a valid identifier and it is not missing or empty. TypeScript expect(order.id).not.toBeNull(); expect(order.id).toBeDefined(); This check is required because the API generates the order ID when a new order is created in the system. It ensures that the “id” field has a valid value generated and assigned to it, since this “id” is used to retrieve, update, or delete order data. TypeScript expect(order.user_id).toEqual("1"); expect(order.product_id).toEqual("79"); expect(order.product_name).toEqual("5 star 10gm Chocobar"); These statements assert that the order details are retrieved correctly for the respective request. The “user_id” - “1” was sent in the request, and verifying it in the response, along with the other order details such as “product_id” and “product_name,” ensures that the correct data is returned. How to Verify Response Data by Matching Objects and Arrays Playwright allows response data verification by matching objects and arrays partially within the API response. This approach is useful because it makes tests more flexible and confirms that the API returns the correct data structure and values. TypeScript test("Get order and verify matching object and array", async ({ request }) => { const response = await request.get("http://localhost:3004/getOrder/", { params: { user_id: 1, }, failOnStatusCode: true, }); const responseBody = await response.json(); expect(responseBody).toMatchObject({ message: "Order found!!", orders: expect.arrayContaining([ expect.objectContaining({ product_id: "79", product_name: "5 star 10gm Chocobar", product_amount: 5, qty: 1, tax_amt: 0.5, total_amt: 5.5, }), ]), }); }); In this test, the toMatchObject assertion verifies that the response contains a “message” with the expected value “Order found!!” and an orders array. Within the array, "expect.arrayContaining" ensures that at least one order matches the expected data, while "expect.objectContaining" verifies only the values in the specified fields of that order. Using Best Practices to Perform Assertions Best practices create stable, maintainable API automation tests by combining basic checks with flexible data matching. TypeScript test("Get Order details API test with best practice", async ({ request }) => { const response = await request.get("http://localhost:3004/getOrder/", { params: { user_id: "1", }, failOnStatusCode: true, }); const responseBody = await response.json(); expect(responseBody.message).toBe("Order found!!"); expect(responseBody.orders.length).toBeGreaterThan(0); expect(responseBody.orders).toEqual( expect.arrayContaining([ expect.objectContaining({ id: 1, product_name: "5 star 10gm Chocobar", }), ]) ); }); The test sends a GET request to fetch order details for “user_id”-“1". The use of failOnStatusCode: true ensures the test fails immediately if the API does not return a 2xx status code. The response is then parsed into a JSON object for validation. The assertions are structured in layers: TypeScript expect(responseBody.message).toBe("Order found!!"); This assertion verifies the message text, confirming that the API returns the correct message when an order is found. TypeScript expect(responseBody.orders.length).toBeGreaterThan(0); This statement ensures meaningful data is returned and avoids false positives when the array is empty. TypeScript expect(responseBody.orders).toEqual( expect.arrayContaining([ expect.objectContaining({ id: 1, product_name: "5 star 10gm Chocobar", }), ]) ); The final part of the code performs the final assertion using arrayContaining and objectContaining to verify that at least one order has the expected “id” and “product_name”, without asserting every field. These layered validations improve clarity by verifying structure, data presence, and key data values in sequence. Extracting Data From the Response Extracting data from the API response is a common and widely used pattern in API test automation. It is important in multiple ways, such as reusing the data in further tests for dynamic testing and end-to-end validation. TypeScript test('Get order details and extract the order id', async({request}) => { const response = await request.get("http://localhost:3004/getOrder/", { params: { id: 1, }, failOnStatusCode: true, }); const responseBody = await response.json(); expect(responseBody.message).toBe("Order found!!"); expect(responseBody.orders.length).toBeGreaterThan(0); expect(responseBody.orders).toEqual( expect.arrayContaining([ expect.objectContaining({ id: 1, product_name: "5 star 10gm Chocobar", }), ]) ); const order = responseBody.orders[0]; expect(order.id).not.toBeNull(); const order_id= order.id; console.log(order_id); const product_name = order.product_name console.log(product_name) }); This test sends a GET API request and performs basic validations to ensure the API response is reliable. TypeScript const order = responseBody.orders[0]; expect(order.id).not.toBeNull(); const order_id= order.id; console.log(order_id); The code above extracts the “order_id” from the order object in the response. Before accessing it, an assertion is made to verify that the value is not null. Finally, the value of the order_id is printed in the console. TypeScript const product_name = order.product_name console.log(product_name) Similarly, other values, such as product_name, can also be extracted. Attaching the Response Body to the Playwright Report The Playwright report, by default, shows the steps executed, the number of tests run, pass/fail status, and time taken to run the tests. However, it does not attach the response body to the test report. Attaching the response body to the report improves visibility and makes the test report more informative and transparent. The following code shows how to extract the required metadata and attach it to the Playwright report. TypeScript test("Get order details API and attach the response details to the report", async ({ request, }, testInfo) => { const response = await request.get("http://localhost:3004/getOrder/", { params: { user_id: "1", }, }); expect(response.status()).toBe(200); const status = response.status(); const statusText = response.statusText(); const headers = response.headers(); const body = await response.json(); const fullResponse = { status, statusText, headers, body, }; await testInfo.attach("Full API Response", { body: JSON.stringify(fullResponse, null, 2), contentType: "application/json", }); }); The testInfo is a built-in Playwright fixture and provides utilities to manage and inspect test execution, such as attaching files to reports, updating test timeouts, and identifying the currently running test. The following lines of code extract the response metadata, such as the status code, status text, headers, and response body. TypeScript const status = response.status(); const statusText = response.statusText(); const headers = response.headers(); const body = await response.json(); Next, let’s combine all response details and create a single object containing: Status codeStatus textHeadersResponse body TypeScript const fullResponse = { status, statusText, headers, body, }; Finally, let’s attach these details to the report using the testInfo.attach() method as shown below: TypeScript await testInfo.attach("Full API Response", { body: JSON.stringify(fullResponse, null, 2), contentType: "application/json", }); The testInfo.attach() adds an attachment to the Playwright report. The attach() method has 3 parameters: Name of the attachment: The first parameter is the name, “Full API Response”, that will be shown for the attachment.Body of the attachment: The second parameter is for the body of the attachment. The JSON.stringify(fullResponse, null, 2) has 3 arguments. The first argument converts the fullResponse object into a readable, pretty-formatted JSON. The second argument is the replacer, which is null. It ensures that all properties from the fullResponse object are included as they are, without modifying anything. The third argument controls pretty-printing. Here, “2” means indent nested JSON by 2 spaces.Content type: This parameter ensures that the report treats the attachment as JSON. The following screenshot is generated after the tests are run: Test Execution Running the tests in Playwright is simple and easy. We can run the following command from the terminal: Plain Text npx playwright test To generate the report, the following command can be used: Plain Text npx playwright show-report Summary Playwright provides multiple approaches, including structure checks and matching objects and arrays for verifying response data. The right strategy should be chosen based on your project’s requirements. Based on my experience, combining response structure checks with response data validation, including the matching object and array strategy, can be used as an effective approach for validating API responses. Happy testing!
Artificial Intelligence has become one of the most influential technologies in modern software development. From chatbots and recommendation systems to sentiment analysis and intelligent search, machine learning models are now expected features in many mobile applications. For several years, integrating AI into iOS applications almost always meant sending user data to cloud services. APIs such as OpenAI, Anthropic Claude, and Google Gemini allowed developers to leverage state-of-the-art language models without worrying about infrastructure or hardware limitations. While this approach is simple, it also introduces several challenges, including network latency, API costs, internet dependency, and privacy concerns. Fortunately, the landscape has changed dramatically. Today's Apple devices contain incredibly powerful hardware, including the Apple Neural Engine (ANE), powerful GPUs, and highly optimized CPUs capable of running sophisticated machine learning models directly on the device. This shift has made on-device AI more practical than ever. Instead of relying entirely on cloud services, developers can now deploy transformer models directly within their applications, enabling offline functionality, lower latency, improved privacy, and reduced operational costs. In this article, we'll explore how to integrate an ONNX-based transformer model into an iOS application using Swift. We'll load a model, execute inference with ONNX Runtime, and prepare the necessary transformer inputs for models such as DistilBERT. A Brief History of LLM Integration in Swift Projects Large Language Models were originally designed to run on powerful cloud infrastructure because of their immense computational requirements. Training these models required thousands of GPUs, and even inference demanded hardware far beyond what smartphones could provide at the time. Because of these limitations, early Swift applications integrated AI almost exclusively through cloud APIs. User prompts were transmitted to remote servers where the model generated a response before sending the results back to the application. Although this architecture worked well, it also introduced unavoidable drawbacks: Internet connectivity became mandatory.Responses depended on network latency.User data had to leave the device.API usage generated recurring operational costs. Meanwhile, Apple continued investing heavily in machine learning acceleration. The introduction of the Apple Neural Engine in 2017 marked a turning point. Every generation of iPhone, iPad, and Mac became increasingly capable of executing neural networks efficiently. At the same time, Apple expanded Core ML, Metal Performance Shaders, and hardware acceleration APIs that allowed developers to run increasingly sophisticated models locally. The open-source AI community accelerated this transition even further. Frameworks such as llama.cpp, MLX, MLC LLM, and ONNX Runtime made it possible to execute optimized transformer models directly on Apple devices. Developers could now deploy popular open-source models including Llama, Mistral, Phi, Gemma, and Qwen without requiring any cloud infrastructure. Apple later introduced Foundation Models as part of Apple Intelligence, further demonstrating the industry's movement toward local AI processing. Today, Swift developers have more choices than ever before. Depending on the application, developers can choose between cloud-hosted models, hybrid cloud/local inference, or fully offline on-device inference. For many applications, including text classification, semantic search, recommendation engines, and lightweight AI assistants, on-device inference has become the preferred solution. Why ONNX? Before diving into the implementation, it's worth understanding why ONNX has become one of the most popular deployment formats for machine learning models. ONNX (Open Neural Network Exchange) is an open standard for representing machine learning models. Instead of locking your project into a specific framework such as TensorFlow or PyTorch, ONNX provides a portable format that can be executed across many different platforms. This portability offers several advantages. A model trained in Python using PyTorch can be exported as an .onnx file and later executed inside an iOS application without rewriting the model itself. Likewise, the exact same model can often be shared between iOS, Android, Windows, Linux, and macOS. This dramatically simplifies deployment across multiple platforms. Microsoft maintains ONNX Runtime, a highly optimized inference engine capable of executing ONNX models efficiently across different hardware accelerators. For Swift developers, this means we only need to load the ONNX model, provide the expected inputs, and retrieve the outputs generated by the runtime. Loading an ONNX Model Let's assume we've already trained our transformer model and exported it to ONNX. Our project now contains a file named MoodClassifier.onnx. The model can either be added directly to the application bundle or packaged as a Swift Package resource. The first step is locating the model inside the application. Swift let modelPath = Bundle.main.path(forResource: "MoodClassifier", ofType: "onnx") If modelPath is not nil, the application has successfully located the model. If it returns nil, verify the following: The model has been added to the target.The filename matches exactly.The resource exists inside the application bundle.The file extension is correct. Successfully locating the model is the first indication that everything has been configured correctly. Installing ONNX Runtime Executing an ONNX model requires an inference engine. Fortunately, Microsoft provides an official Swift Package for ONNX Runtime that can be added using Swift Package Manager. Swift .package(url: "https://github.com/microsoft/onnxruntime-swift-package-manager",from: "1.24.2") After adding the dependency, we're ready to create an inference session. Running Inference Running inference simply means executing a trained machine learning model using new input data. Unlike training, inference does not modify the model. It only computes predictions. Creating an ONNX Runtime session requires three primary components: ORTEnvORTSessionOptionsORTSession The environment configures runtime behavior and logging. The session options allow developers to customize execution behavior. Finally, the session loads the model into memory and prepares it for inference. Swift let env = try ORTEnv(loggingLevel: .warning) let options = try ORTSessionOptions() let modelPath = Bundle.main.path(forResource: "MoodClassifier", ofType: "onnx")! let session = try ORTSession(env: env, modelPath: modelPath, sessionOptions: options) let outputs = try session.run(withInputs: [:], outputNames: ["logits"], runOptions: nil) In this example, the model returns a tensor called logits. Depending on how the model was exported, your output tensor may have a different name. Always inspect the exported model to determine the available output names. Preparing Inputs for Transformer Models Most transformer models, including DistilBERT, BERT, and RoBERTa, cannot process raw text directly. Instead, they expect numerical tensors representing the input sentence. This process is called tokenization. Tokenization converts natural language into token IDs that correspond to entries within the model's vocabulary. Alongside the token IDs, transformer models also require an attention mask. The attention mask tells the model which tokens belong to the original sentence and which tokens are merely padding added to maintain a fixed sequence length. Using the correct tokenizer is extremely important. The tokenizer used during inference must be identical to the tokenizer used while training the model. Even small differences in vocabulary or preprocessing rules can generate completely different token IDs, resulting in poor predictions despite using the correct model. Using Swift Transformers Hugging Face provides an excellent package called Swift Transformers that simplifies tokenization directly within Swift. The package can be added using Swift Package Manager. Swift .package(url: "https://github.com/huggingface/swift-transformers", from: "1.3.3" ) Once installed, you can load the tokenizer that matches the model used during training and generate the input_ids and attention_mask required by the transformer. After generating these arrays, they must be converted into ONNX tensors before inference. Creating ONNX Input Tensors The generated token arrays must be wrapped inside ORTValue tensors. The following example converts both arrays into tensors before executing the model. Swift let env = try ORTEnv(loggingLevel: .warning) let options = try ORTSessionOptions() let modelPath = Bundle.main.path(forResource: "MoodClassifier", ofType: "onnx")! let session = try ORTSession(env: env, modelPath: modelPath, sessionOptions: options) let inputData = NSMutableData(bytes: &inputIDs, length: inputIDs.count * MemoryLayout<Int64>.size) let inputTensor = try ORTValue(tensorData: inputData, elementType: .int64, shape: [1, inputIDs.count] as [NSNumber]) let attentionData = NSMutableData(bytes: &attentionMask, length: attentionMask.count * MemoryLayout<Int64>.size) let attentionTensor = try ORTValue(tensorData: attentionData, elementType: .int64, shape: [1, attentionMask.count] as [NSNumber]) let outputs = try session.run(withInputs: ["input_ids": inputTensor, "attention_mask": attentionTensor], outputNames: ["logits"], runOptions: nil) In this example, two tensors are created: input_ids, which contains the numerical representation of the input text, and attention_mask, which tells the model which tokens should participate in the attention mechanism. These tensors are then passed into the ONNX Runtime session, which executes the model and returns the requested outputs. Conclusion On-device AI is no longer a niche capability reserved for flagship applications. Thanks to frameworks such as ONNX Runtime and Swift Transformers, integrating transformer models into iOS projects has become both accessible and practical. In this article, we explored how to load an ONNX model, execute inference using Microsoft's ONNX Runtime, and prepare the input_ids and attention_mask tensors required by transformer models. These components form the foundation for deploying a wide range of AI-powered features directly within Swift applications. As Apple's hardware continues to evolve and transformer models become increasingly efficient, local inference will play an even greater role in the future of mobile development. Whether you're building a sentiment analyzer, semantic search engine, recommendation system, or lightweight AI assistant, ONNX Runtime provides a robust and portable solution for bringing modern machine learning to iOS.
Before diving into solutions, it helps to understand the scale of the problem. Take a common production pattern: a customer support bot that processes 10,000 messages per day, each with a 2,000-token system prompt and a 200-token user message. Cost breakdown by the numbers: scenariomodeldaily costNo optimizationClaude Opus ($5/1M)$11.00Right model for taskClaude Haiku ($1/1M)$2.20Add prompt cachingHaiku + caching ($0.10/1M cached)$0.42Combined savings96% reduction That's not a benchmark; it's arithmetic. The techniques don't require magic; they require applying what providers already offer. Technique 1: Prompt Caching, Up to 90% Off Repeated Content The Problem Most LLM applications send the same system prompt on every request. If your system prompt is 2,000 tokens and you make 10,000 requests per day, you're paying for 20 million input tokens daily, even though the content never changes. How It Works Anthropic's prompt caching lets you mark content blocks with cache_control. The first request pays full price and writes to cache. Every subsequent request that hits the same cached content pays 10% of the normal input price. Cache entries last 5 minutes and reset on each hit. The key insight: cache the most stable content first. Your base instructions change rarely. Your few-shot examples change occasionally. Your per-request context changes every time. Structure your prompt from most stable to least stable. Python from llm_optimizer import OptimizedClient, build_cached_system_prompt import anthropic client = OptimizedClient(anthropic_client=anthropic.Anthropic()) # Build an optimally structured cached system prompt system = client.build_cached_system( base_instructions=""" You are an expert customer support agent for a SaaS company. You have deep knowledge of our product, billing, and technical issues. Always be empathetic, clear, and solution-focused. [... 1,500 more tokens of stable instructions ...] """, # ← cached after first request — 10% cost on all subsequent calls few_shot_examples=""" Example 1: Billing question → here's how to handle it Example 2: Technical issue → here's the escalation path [... 500 tokens of examples ...] """, # ← also cached separately ) # First call: pays full price, writes to cache response1 = client.complete(messages=[{"role": "user", "content": "How do I cancel?"}], system=system) # Second+ calls: system prompt served from cache at 10% cost response2 = client.complete(messages=[{"role": "user", "content": "Where's my invoice?"}], system=system) What the Library Does llm-optimizer automatically injects cache_control breakpoints at optimal positions, system prompt, few-shot examples, and long conversation history, respecting Anthropic's 4-breakpoint limit. You don't touch the API directly. Savings Calculation Shell 2,000 token system prompt × 10,000 requests/day = 20M tokens/day Without caching: 20M × $3.00/1M (Sonnet) = $60.00/day With caching: 2M × $3.00 + 18M × $0.30 = $11.40/day Savings: $48.60/day = $17,739/year Technique 2: Model Routing, 60% to 80% Off by Using the Right Model The Problem Routing every request to your best model is the most common and most expensive mistake. Claude Opus costs 5x more than Claude Haiku. For tasks that Haiku handles perfectly, such as classification, extraction, translation, and simple Q&A, you're paying a 500% premium for no benefit. The Naive Approach and Why It Fails The obvious solution is to route by keyword: if the prompt contains "classify," use Haiku; if it contains "analyze," use Sonnet. This works until it doesn't. A prompt like "Explain the constitutional implications of this clause" is 8 words. Short, simple-looking. A keyword router sees no complexity signals and routes it to Haiku. But the task requires expert-level legal reasoning. This class of error intent-heavy short prompts is the primary failure mode of heuristic routing. The Better Approach: Use a Classifier llm-optimizer solves this by optionally using Haiku itself to classify task complexity before routing. The cost is approximately 15 tokens, about $0.000015. If that classification prevents one wrong Opus call (2,000 tokens × $5/1M = $0.01), it pays for itself 666 times over. Python from llm_optimizer import OptimizedClient, Provider client = OptimizedClient( anthropic_client=anthropic.Anthropic(), enable_llm_classifier=True, # uses Haiku to assess complexity — ~$0.000015/call preferred_provider=Provider.ANTHROPIC, ) # Short prompt, complex intent → correctly routed to Opus response = client.complete( messages=[{"role": "user", "content": "Explain the constitutional implications of this clause"}] ) # You can audit the routing decision from llm_optimizer import ModelRouter router = ModelRouter(enable_llm_classifier=False) print(router.explain("classify this email as spam or not")) # { # "detected_complexity": "simple", # "routed_model": "claude-haiku-4-5", # "keyword_signals_fired": {"simple": ["classify"]}, # "token_count": 8 # } Complexity Tiers A closer look at complexity: Tierexamplesdefault modelSimpleClassification, extraction, yes/no, translationClaude HaikuMediumSummarization, paraphrasing, short Q&AClaude HaikuComplexCode generation, analysis, evaluationClaude SonnetExpertLegal reasoning, research, system design, math proofsClaude Opus Shell 1,000 requests/day — mixed complexity Without routing: all → Opus ($5/1M input) 1,000 × 500 tokens = 500K tokens × $5 = $2.50/day With routing: 70% Haiku, 20% Sonnet, 10% Opus 700 × 500 × $1 + 200 × 500 × $3 + 100 × 500 × $5 = $1.05/day Savings: 58% reduction Technique 3: Prompt Optimization, 5% to 20% Off Token Count The Problem Prompts written by humans, especially in collaborative or enterprise settings, accumulate filler. Phrases like "please note that", "it is important to note that", "in order to", and "due to the fact that" add tokens without adding meaning. At scale, this is a measurable cost. What the Library Strips Python from llm_optimizer import PromptOptimizer opt = PromptOptimizer() result = opt.optimize(""" In order to complete this task, please note that you should carefully analyze the following text. It is important to note that accuracy matters. Please be aware that your response should be concise. """) print(result.optimized_text) # "To complete this task, carefully analyze the following text. # Accuracy matters. Your response should be concise." print(f"Saved {result.tokens_saved} tokens ({result.savings_pct}%)") # Saved 18 tokens (31%) What is never touched: Code blocks, factual content, user-specified phrasing. The optimizer is conservative by default. It only removes patterns with no semantic value. Conversation history trimming: In long conversations, the library keeps the last N turns and drops older context, preventing unbounded token grow. Technique 4: Batch Processing, 50% Off Non-Urgent Requests The Problem Not every LLM call needs an immediate response. Nightly report generation, document indexing, data enrichment pipelines, and offline classification jobs all of these run fine with a delay. But most teams send them as real-time requests anyway, paying full price. How Anthropic's Batch API Works Anthropic's Message Batch API processes up to 10,000 requests per batch at 50% of the normal price. Results are available within minutes to hours. The trade-off is explicit: cost for latency. Python client = OptimizedClient( anthropic_client=anthropic.Anthropic(), enable_batching=True, ) # Queue 1,000 document summaries throughout the day for doc in documents: client.queue( custom_id=doc["id"], messages=[{"role": "user", "content": f"Summarize: {doc['text']}"}], max_tokens=200, ) # Submit as one batch — 50% cheaper than 1,000 individual calls batch_id = client.submit_batch() # Poll when ready — minutes to hours depending on load results = client.poll_batch(batch_id, wait=True) for r in results: print(f"{r.custom_id}: {r.content}") When to Use It Nightly data processing pipelinesDocument indexing and enrichment Offline classification and tagging Report generation Technique 5: Document Compression to Reduce Context Before Sending The Honest Tradeoff This technique requires a direct warning: document compression is lossy. Removing content from a document to reduce token count means the model works with less information. For some tasks this is fine; for others it produces wrong answers. Use it only when: You've verified empirically that compression doesn't hurt your answer quality.You're doing rough extraction where completeness isn't required.You have re-ranking downstream (e.g., RAG pipelines). Do not use it for legal documents, compliance reviews, or any task where every sentence may be relevant. TF-IDF Extractive Compression When you do use compression, naive truncation (cutting from the end) is the worst strategy. llm-optimizer implements TF-IDF paragraph scoring where each paragraph is scored by its term overlap with your query, weighted by how unique those terms are across the document. The most relevant paragraphs fill the token budget; the rest are dropped. Python from llm_optimizer import DocumentCompressor # ⚠️ Read the accuracy warning before using in production comp = DocumentCompressor( max_tokens=4000, strategy="extractive", # TF-IDF scoring — best accuracy ) compressed, tokens_saved = comp.compress( document=long_contract, # 50,000 tokens query="payment terms and termination clauses" # focus compression here ) print(f"Compressed to {4000} tokens, saved {tokens_saved} tokens") # All compressed output includes a visible [⚠️ COMPRESSION WARNING] marker There are three strategies available: Strategy options: StrategyHow it worksbest forExtractiveTF-IDF scoring against queryWhen you have a specific querySmartKeeps first 60% + last 20%Structured documents with summariesTruncateHard cutoffWhen you need predictable behavior Technique 6: Cost Tracking That Observes Before You Optimize Why This Matters You can't optimize what you don't measure. Before applying any of the above techniques, you need to know: Which models you're actually usingWhere your token spend is goingWhether your optimizations are working Python client = OptimizedClient( anthropic_client=anthropic.Anthropic(), persist_tracking="usage.jsonl", # survives restarts ) # ... run your application ... client.print_summary() # ═══════════════════════════════════════════════════════ # LLM Cost Optimizer — Usage Summary # ═══════════════════════════════════════════════════════ # Total Requests : 1,247 # Total Cost : $0.8432 # Total Saved : $7.2180 (89.5% savings) # Cached Tokens : 8,432,000 # By Model : haiku: 891 reqs ($0.12) | sonnet: 312 ($0.58) # Optimizations : prompt_caching: 1247x | model_routing: 1247x # ═══════════════════════════════════════════════════════ The tracker records every request's tokens, cost, cached tokens, savings, latency, and which optimizations fired. Data persists to JSONL so you can analyze it across sessions or pipe it to your observability stack. Architecture: Why Not Just Use LiteLLM? The obvious question. LiteLLM is excellent and covers a lot of ground, including unified provider API, routing, cost tracking, batch processing. If you're not already using it, you should evaluate it. llm-optimizer does three things LiteLLM doesn't: Automatic cache_control injection: LiteLLM passes caching headers through but doesn't inject breakpoints at optimal positions automatically.Prompt filler stripping: LiteLLM has no token-level prompt optimization.TF-IDF document compression: LiteLLM has no query-aware document compression. The intended use is actually as a complement: llm-optimizer can sit on top of a LiteLLM setup, handling the prompt-level optimizations that LiteLLM doesn't touch. Streaming Support For user-facing applications, the library supports streaming: Python with client.stream( messages=[{"role": "user", "content": "Explain quantum entanglement"}], system="You are a physics tutor.", max_tokens=512, ) as stream: for chunk in stream: print(chunk, end="", flush=True) # Access token usage after stream completes usage = stream.usage() All optimizations, including caching, routing, prompt, and optimization, apply identically to streaming requests. Error Handling Production LLM applications need to handle rate limits and model overloads gracefully. The library handles this automatically: Python client = OptimizedClient( anthropic_client=anthropic.Anthropic(), max_retries=3, # retry on rate limit with exponential backoff retry_base_delay=1.0, # 1s, 2s, 4s ) # Rate limit (429) → retried with backoff # Model overloaded (529) → falls back to next capable model automatically # Non-retriable error → raises immediately Limitations Honest about what this doesn't do yet: Not org-scale validated: v0.4.0 is tested against ai-core Bedrock and Anthropic direct (14 live tests, all passing). Not yet run against production workloads at team scale. The pilot measures this.No async support: complete() and stream() are synchronous. Async support planned for a future release.OpenAI and Google partially tested: Anthropic and AWS Bedrock are the validated providers. OpenAI is implemented but not end-to-end tested in CI. Google Gemini streaming is not yet implemented.Token counting is approximate: The default estimator is within ~20% of the actual count. Install tiktoken for exact counts: pip install llm-optimizer[tiktoken].Pricing data can go stale: Stored in pricing.json with a version stamp. The library warns automatically if data is older than 30 days.Model allowlist: Only models listed in pricing.json can be routed to. Mythos and Fable 5 are not in the registry and cannot be called. Adding a new model requires a deliberate update to pricing.json. Python # Basic install pip install llm-optimizer # With exact token counting pip install llm-optimizer[tiktoken] # All providers pip install llm-optimizer[all] Python import anthropic from llm_optimizer import OptimizedClient client = OptimizedClient( anthropic_client=anthropic.Anthropic(), # All optimizations on by default except compression (lossy — opt-in) ) response = client.complete( messages=[{"role": "user", "content": "Your prompt here"}], system="Your system prompt here", ) client.print_summary() Links: PyPI: https://pypi.org/project/llm-optimizeGitHub: https://github.com/banerjeeso/llm-optimiz What's Next Async support (acomplete(), astream())LiteLLM adapterBudget guard: raise before a request exceeds a cost thresholdReal production benchmarks once I've run this against a live workload. Feedback welcome, especially from anyone who works with LLM APIs in production and can stress-test the routing logic or compression accuracy. Published on PyPI as llm-optimizer. MIT license. Contributions welcome.
Artificial intelligence is transforming software testing by enabling faster test creation, smarter execution, and more efficient quality assurance processes. With AI agents, we can generate test cases and scripts, run tests, and produce detailed reports with minimal manual effort. AI agents are software systems that leverage artificial intelligence to achieve goals and perform tasks on behalf of users. They can think through problems, plan actions, and remember things, while also making decisions on their own and improving over time. In this tutorial, we will build an AI-powered agent that does exactly this. It takes a .txt file with a simple use case written in plain English, processes it using an AI model such as OpenAI or Ollama, and generates Selenium test automation scripts in Java. How to Build an AI Agent for Test Automation We will be building an AI agent for test automation that will perform the following tasks: Reads a test scenario from a .txt fileSends the scenario with the right prompt to OpenAI/Ollama, based on the provider selected in the config fileGenerates the Selenium WebDriver Java test automation scripts bifurcated into: Page Object ClassesTest ClassTestng.xmlReadme.md with notes and steps to run the test Prerequisites The following prerequisites are required to start building the AI agent for test automation: Python 3.14 and aboveCode Editor (VS Code is preferred)Basic understanding of Python and the TerminalTest scenario written in plain English text LLM Setup OpenAI API KeyOllama setup on the local machine Understanding the Selenium AI Test Generation Workflow Before we deep dive into building the AI agent, let’s get a high-level overview of how this AI Agent works for generating Selenium Java test automation scripts: The process is explained step by step below: User input: The user provides a plain English test scenario in a .txt file.Processing layer: The Python main script reads the test scenario file and forwards it to the LLM model (OpenAI or Ollama), depending on the configuration.AI engine: Before sending requests to the LLM, the AI agent converts the test scenario from plain English to a structured prompt and adds clear instructions like: Generate Selenium WebDriver code in JavaUse the TestNG frameworkFollow the Page Object ModelGenerate Multiple Page Object classes(if needed)Rules for generating the codeOutput format for files The LLM understands the test description, steps, and expected behavior, and converts them into the respective files.File generation: The AI agent creates a new timestamp-based folder and stores the generated code in separate files, bifurcating the page object classes, test class, testng.xml, and README.md. Markdown . ├── config.py ├── tools.py ├── main.py ├── requirements.txt ├── .env ├── test_cases/ │ ├── input/ │ │ └── sample_test_case.txt │ └── output/ Let’s start building the AI agent using the following step-by-step guide: Step 1: Creating and Activating a Python Virtual Environment Run the following command in the terminal to create a Python virtual environment: Plain Text python -m venv venv Next, run the following command to activate the virtual environment: On MacOS/Linux: Plain Text source venv/bin/activate On Windows: Plain Text venv\Scripts\activate Step 2: Installing Dependencies Let’s create a requirements.txt as it serves as a blueprint for recreating the exact environment needed to run the project consistently across different environments. Plain Text #requirements.txt openai==2.28.0 requests==2.32.5 python-dotenv==1.2.2 The following command should be run from the terminal to install the dependencies: Plain Text pip install -r requirements.txt Step 3: Writing the Configuration File The configuration file holds details related to OpenAI, Ollama, input and output files, and settings for using the desired LLM provider to generate test automation scripts. Python #config.py from dataclasses import dataclass, field from pathlib import Path from dotenv import load_dotenv from dotenv import load_dotenv from openai import OpenAI import os load_dotenv() @dataclass class OpenAIConfig: model_name: str = os.getenv("OPENAI_MODEL_NAME", "gpt-5") max_tokens: int = int(os.getenv("OPENAI_MAX_TOKENS", "4096")) temperature: float = float(os.getenv("OPENAI_TEMPERATURE", "0.3")) api_key: str = os.getenv("OPENAI_API_KEY", "") @property def client(self): if not self.api_key: raise ValueError( "OPENAI_API_KEY is required when using OpenAI provider" ) return OpenAI(api_key=self.api_key) @dataclass class OllamaConfig: model: str = os.getenv("OLLAMA_MODEL_NAME", "llama3") prompt: str = "You are a Selenium WebDriver test automation expert. Generate Selenium WebDriver Java test scripts. STRICTLY follow the format. Any deviation is not acceptable." stream: bool = False temperature: float = float(os.getenv("OLLAMA_TEMPERATURE", "0.3")) ollama_baseurl:str = os.getenv("OLLAMA_BASEURL", "http://localhost:11434") ollama_endpoint: str = os.getenv("OLLAMA_ENDPOINT", "http://localhost:11434/api/generate") @dataclass class FileConfig: input_file: Path = Path("test_cases/input/sample_test_case.txt") output_file_path: Path = Path("test_cases/output/") @dataclass class AppConfig: provider:str = "ollama" #openai or ollama openai: OpenAIConfig = field(default_factory=OpenAIConfig) ollama: OllamaConfig = field(default_factory=OllamaConfig) files: FileConfig = field(default_factory=FileConfig) config = AppConfig() The configuration file serves as a centralized configuration system for switching between LLM providers and managing file settings. Let’s break it down to understand it further: Dataclasses: The @dataclass automatically generates the __init__() method, making it easier to create and manage objects without writing boilerplate code.Separate Configs: Separate config classes are defined for OpenAI, Ollama, and input/output file handling, each with default values. The input file name is currently configured as “sample_test_case.txt.”AppConfig: The AppConfig class serves as a central wrapper, enabling switching between providers (OpenAI or Ollama) with a single flag.field(default_factory=…): It ensures each config gets its own instance, avoiding shared state issues. Step 4: Implementing Utility Functions for Reading and Generating Code In this step, we’ll create a new Python file named tools.py and implement the following utility functions within it: load_test_case_from_file()build_prompt()generate_with_openai()generate_with_ollama()generate_selenium_test_script()create_timestamped_output_dir()split_and_save_files() The following packages should be imported in the tools.py file: Python import requests from config import config from typing import Optional from pathlib import Path from datetime import datetime from logger import logger import requests from config import config from logger import logger Let’s learn the utility functions one by one. load_test_case_from_file() function: Python def load_test_case_from_file(file_path: str | Path) -> str: logger.info(f"Loading test case from: {file_path}") try: with open(file_path, "r", encoding="utf-8") as file: content = file.read() logger.info(f"Test case loaded successfully ({len(content)} characters)") return content except FileNotFoundError: logger.error(f"Test case file not found: {file_path}") raise FileNotFoundError(f"Test case file not found: {file_path}") except Exception as e: logger.exception(f"Error reading test case file: {e}") raise RuntimeError(f"Error reading test case file: {e}") The load_test_case_from_file() function takes the test case file path as a parameter, reads its contents, and returns them as a string. It safely handles errors by raising a clear message if the file is not found and wraps any other unexpected issues in a RuntimeError. build_prompt() function: Python def build_prompt(use_case_text: str) -> str: logger.info("Building AI prompt") prompt = f""" You are a test automation expert specializing in Selenium WebDriver with Java. Generate Selenium automation test script using the following instructions and skills: - Use Java 17 to write the code - Use latest Selenium WebDriver Java dependency version to write code - Do not write code statement "System.setProperty()" to add chromedriver path, in the tests - Follow Page Object Model (POM) - Use latest version of TestNG dependency - Apply best coding practices for writing Java code - Add comments explaining each step - Add assertions using TestNG assertion - Do not add random assertion statements in the code - Do not mention any text in README that says that ChromeDriver path should be added to the Path - Never use brittle XPATH and CSS Selectors selectors such as .btn-primary, .container > div:nth-child(2), #content div span, or auto-generated classes. IMPORTANT: You MUST follow the exact output format below. Rules: - ALWAYS start each file with ===FILE: filename=== - Use class names based on the web page(e.g. HomPage.java, LoginPage.java, etc. These names are for instructions only, use class name specific to the web page) - Do not add "Page" to the test class name - Do NOT add explanations outside file blocks - DO NOT skip this format - If you do not follow this format, the output will be rejected - DO NOT use markdown (no **, no ``` blocks) - DO NOT add file names outside ===FILE: markers - ONLY use ===FILE: filename=== format - OUTPUT FORMAT FOR FILES(STRICT): ===FILE: filename=== file content - The following files MUST only be generated in the same order(STRICT). No deviation is acceptable: - Multiple Page Object classes(if needed) (Strictly Page object class, no WebDriver instantiation in these classes, Do not create duplicate page object classes) - Test class(WebDriver should be instantiated in the Test class, Do not use WebDriverManager to instantiate WebDriver, Use TestNG's @BeforeMethod annotation and define a method to instantiate the WebDriver, Use TestNG's @AfterMethod to quit the WebDriver) - Add assertions using TestNG assertion - Do not add random assertion statements in the code - testng.xml(Follow correct structure as per TestNG guidelines) - README.md (Include notes and steps to run the test using testng.xml file) Use Case: {use_case_text} """ logger.info(f"Prompt created ({len(prompt)} characters)") return prompt The build_prompt() function is a core part of the application because it constructs the prompt sent to the LLM to generate Selenium Java test automation code. It takes the use case text as input and embeds it into a detailed instruction template, guiding the model to generate Java-based Selenium tests using best practices such as POM and TestNG. The prompt also enforces file-formatting rules like (===FILE: filename===) to ensure the output is structured, consistent, and ready for file generation without manual cleanup. It also enforces rules to generate the POM, test class, testng.xml, and README in strict order to keep the output consistent. Using Effective Prompts Effective prompts are essential for generating clean, reliable Selenium test scripts with AI. By clearly defining the requirements, the model can be guided to use best coding practices, create clean test scripts, and use the Page Object Model (POM) to generate well-structured code. Prompt Design for Better Selenium Java Tests Clear, specific prompts help the AI generate cleaner Selenium test scripts. Mention details such as using Java for Selenium, using the latest versions of Selenium and TestNG, and following coding best practices so the AI can better meet the requirements and produce well-designed output. Customizing Prompts for Page Object Model (POM) The following prompts can be used to train the AI to generate the page object files we need : Use best practices to generate the Page Object classes.Generate Multiple Page Object classes (if needed).Do not instantiate the WebDriver in the Page Object class.Do not create duplicate Page Object classes.Additionally, we can also mention using best locator strategies, like preferring IDs and CSS selectors over complex XPaths, to make the tests more stable, readable, and easy to maintain. generate_with_openai() function: Python def generate_with_openai(prompt: str) -> Optional[str]: response = config.openai.client.chat.completions.create( model=config.openai.model_name, temperature=config.openai.temperature, max_tokens=config.openai.max_tokens, messages=[ {"role": "system", "content": "You are a Selenium WebDriver test automation expert. Generate Selenium WebDriver Java test scripts. STRICTLY follow the format. Any deviation is not acceptable."}, {"role": "user", "content": prompt}, ], ) return response.choices[0].message.content The generate_with_openai() function sends a prompt to the OpenAI API to generate Selenium WebDriver test scripts in Java. The OpenAI client is initialized using the API key from an environment variable, allowing the code to connect and interact with OpenAI services securely. It uses the model, temperature, and token limits from the configuration. It generates the request with a message that defines the AI’s role and user prompt (the build_prompt() function will be supplied here). The API returns multiple choices, and the function extracts the generated content from the first response. Finally, it returns the generated test script as a string. generate_with_ollama() function: Python def generate_with_ollama(prompt: str) -> str: logger.info( f"Generating test code with Ollama model: {config.ollama.model}" ) response = requests.post( config.ollama.ollama_endpoint, json={ "model": config.ollama.model, "prompt": config.ollama.prompt + "\n" + prompt, "stream": config.ollama.stream, }, ) response.raise_for_status() data = response.json() generated_text = data.get("response", "") logger.info( f"Ollama generation completed ({len(generated_text)} characters)" ) return generated_text The generate_with_ollama() function generates text by sending a POST request to the configured Ollama API endpoint. It combines a predefined base prompt with the user-provided prompt and sends it along with the model name and streaming configuration. After receiving the response, raise_for_status() checks for HTTP errors, and response.json() parses the response into a Python dictionary. The generated text is extracted from the response field, with an empty string used if the field is missing. Finally, the function logs the response length and returns the generated text. generate_selenium_test_script() function: Python def generate_selenium_test_script(test_case_text: str) -> Optional[str]: logger.info("Starting Selenium test generation") prompt = build_prompt(test_case_text) provider = config.provider.lower() logger.info(f"Using AI provider: {provider}") if provider == "openai": logger.info("Sending prompt to OpenAI") return generate_with_openai(prompt) elif provider == "ollama": logger.info("Sending prompt to Ollama") return generate_with_ollama(prompt) else: logger.error(f"Unsupported provider: {provider}") raise ValueError(f"Unsupported provider: {provider}") The generate_selenium_test_script() function acts as a wrapper to generate a Selenium WebDriver test script based on the input test case. It first builds a prompt using the build_prompt (test_case_text) function, which prepares the input for the LLM. Next, based on the configured provider, i.e., OpenAI or Ollama, it dynamically calls the respective function to generate the script. If an unsupported provider is specified, it raises a ValueError to prevent unexpected behavior. create_timestamped_output_dir() function: Python def create_timestamped_output_dir(base_output_path: Path) -> Path: timestamp = datetime.now().strftime("%Y-%m-%d_%H-%M-%S") output_dir = base_output_path / timestamp output_dir.mkdir(parents=True, exist_ok=True) logger.info(f"Created output directory: {output_dir}") return output_dir The create_timestamped_output_dir() function creates a new folder named with the current date and time, so each run gets its own unique output directory. It ensures the folder exists (creating it if needed) and returns its path for saving files. split_and_save_files() function: Python def split_and_save_files(generated_text: str, base_output_path: Path) -> None: sections = generated_text.split("===FILE:") if len(sections)<=1: logger.error("No structured files found in AI response") raise ValueError ("No Structured files found in AI response!") logger.info(f"Found {len(sections) - 1} generated files") pageobject_dir = base_output_path/"pageobjects" pageobject_dir.mkdir(parents=True, exist_ok=True) for section in sections[1:]: section = section.strip() parts = section.split("\n",1) raw_filename = parts[0].strip() filename = raw_filename.replace("===", "").strip().split()[0] content = parts[1].strip() if len(parts)>1 else "" content = content.split("===FILE:")[0].strip() if filename.endswith("Page.java"): file_path = pageobject_dir / filename else: file_path = base_output_path / filename logger.info(f"Writing file: {file_path}") with open(file_path, "w", encoding="utf-8") as f: f.write(content) logger.info(f"Created: {file_path} ({len(content)} characters)") The split_and_save_files() function takes the AI-generated response and splits it into multiple files based on a marker (===FILE:). It extracts each filename and its content, then saves them into the appropriate folders. The Page Object files go into a pageobjects directory, while the test class, testng.xml, and README.md files go into the main output/<timestamped> folder. If the expected structure is missing, it throws an error to avoid saving incorrect output. Step 5: Writing the Main Script to Glue Everything In this step, we’ll create a main.py file to orchestrate the complete workflow of generating Selenium test scripts using AI: Python import re from tools import load_test_case_from_file, generate_selenium_test_script,split_and_save_files,create_timestamped_output_dir, check_llm_connection from config import config from pathlib import Path from logger import logger def main() -> None: try: logger.info("Starting Selenium AI Test Generator") if not check_llm_connection(): logger.error( "❌ LLM connection check failed. " "Stopping test generation." ) return input_file = config.files.input_file output_file_path = config.files.output_file_path run_output_dir = create_timestamped_output_dir(output_file_path) use_case_text = load_test_case_from_file(input_file) if not use_case_text: logger.error("Test case file is empty.") raise ValueError("Test case file is empty.") generated_output = generate_selenium_test_script(use_case_text) if not generated_output: logger.error("Failed to generate Selenium WebDriver Java test automation scripts.") raise RuntimeError("Failed to generate Selenium WebDriver Java test automation scripts.") split_and_save_files(generated_output,run_output_dir) logger.info(f"✅ All generated files saved successfully to {run_output_dir}") except Exception as e: logger.error(f"❌ Error occurred while generating output files: {e}") if __name__ == "__main__": main() The main.py file is where the program starts and manages the whole process of generating Selenium test scripts from start to finish. It reads the input test case file, sends the test case to the AI to generate automation scripts, and generates the timestamped output folder. Once the scripts are generated, it splits them into multiple files and saves them in the appropriate directories. In case of any exception, it catches the error and prints the message “Error occurred while generating output files” with the exception details. Generating First Test Scripts Let’s create a new text file, ”sample_test_case.txt,” and place it in the input/ folder with the following test scenario to generate the Selenium test scripts using the AI agent we created: Plain Text Title: Application Login scenario Precondition: User is registered in the application. Steps: 1. Open Chrome browser 2. Navigate to https://ecommerce-playground.lambdatest.io/index.php?route=account/login 3. Enter "[email protected]" in the E-Mail Address field 4. Enter "Password@321" in the Password field 5. Click on the Login Button 5. Add an assert statement to check that "My Account" page is displayed. Configuration and Setup Using OpenAI Create a .env file and add your OpenAI API key in it. This file should always remain on your local machine and should not be committed to the remote repository. Plain Text OPENAI_MODEL_NAME=<model name> OPENAI_MAX_TOKENS=<max tokens> OPENAI_TEMPERATURE=<temperature value> OPENAI_API_KEY=<Your OpenAI API Key> Using Ollama Download and install Ollama.Start Ollama by running the command ollama serve in the terminal.Pull the Qwen3:8b model by running the command — ollama pull qwen3:8b. Create a .env file and update the following details in it: Plain Text OLLAMA_MODEL_NAME=qwen3:8b OLLAMA_TEMPERATURE=0.3 OLLAMA_ENDPOINT=http://localhost:11434/api/generate OLLAMA_BASEURL=http://localhost:11434 Before running the model, ensure that Ollama is running in the background. Open Terminal and run the command ollama serve Note: Running the model explicitly is not required. Running the AI Agent Open the terminal and run the following command to generate the test scripts: Plain Text python main.py The following log should be printed in the console after the Agent is run successfully: The test scripts should be generated in a new timestamped(YYYY-MM-DD_hh-mm-ss) folder that is generated inside the output/ folder: Reviewing the Generated Java Selenium Code Let’s check each file generated in the output/ folder and review the code to verify if it fits the test scenario we defined and follows the expected structure and best practices. Let’s check each file generated in the output/ folder and review the code to verify if it fits the test scenario we defined and follows the expected structure and best practices. The following page object class is generated: LoginPage.java Java public class LoginPage { private WebDriver driver; public LoginPage(WebDriver driver) { this.driver = driver; } public void navigateToLoginPage() { driver.get("https://ecommerce-playground.lambdatest.io/index.php?route=account/login"); } public void enterEmail(String email) { driver.findElement(By.name("email")).sendKeys(email); } public void enterPassword(String password) { driver.findElement(By.name("password")).sendKeys(password); } public void clickLoginButton() { driver.findElement(By.xpath("//button[@type='submit']")).click(); } } The page object file is generated correctly using the name and XPath locator strategies, and appropriate methods are created to interact with the respective WebElements on the page. However, the import statements are missing, which should be added to avoid errors. Additionally, locators should be checked, and explicit waits could be added in this class while locating the elements to reduce flakiness during test execution. The following test class is generated: LoginTest.java Java import org.openqa.selenium.By; import org.openqa.selenium.WebDriver; import org.openqa.selenium.WebElement; import org.openqa.selenium.chrome.ChromeDriver; import org.testng.Assert; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Test; public class LoginTest { private WebDriver driver; @BeforeMethod public void setup() { System.setProperty("webdriver.chrome.driver", "path_to_chrome_driver"); driver = new ChromeDriver(); } @AfterMethod public void teardown() { driver.quit(); } @Test public void testLoginScenario() { LoginPage loginPage = new LoginPage(driver); loginPage.navigateToLoginPage(); loginPage.enterEmail("[email protected]"); loginPage.enterPassword("Password@321"); loginPage.clickLoginButton(); WebElement myAccountPageTitle = driver.findElement(By.xpath("//h1[@class='page-title']")); Assert.assertTrue(myAccountPageTitle.isDisplayed(), "My Account page is not displayed"); } } The test class is generated per the provided prompt and includes the @BeforeMethod and @AfterMethod annotations for setup and teardown. However, it includes System.setProperty(), which can be removed because the latest version of Selenium doesn't require it to set the ChromeDriver path. The import statement for the @AfterMethod annotation is missing, which should be added. Additionally, the myAccountPageTitle WebElement can be moved to a separate MyAccount Page Object class, and its XPath should also be verified to ensure it's pointing to the correct element on the page. Testng.xml XML <?xml version="1.0" encoding="UTF-8"?> <!DOCTYPE suite SYSTEM "http://testng.org/testng-1.0.dtd"> <suite name="LoginTestSuite"> <test name="LoginTest"> <classes> <class name="LoginTest"/> </classes> </test> </suite> The qualified class name should be updated once the code is moved to the project. README.md Markdown **How to run the test** 1. Install TestNG framework and Selenium WebDriver. 2. Place the Selenium ChromeDriver executable in a directory, update the path_to_chrome_driver variable in `LoginTest.java` with this directory. 3. Run the test using the testng.xml file by executing the command: `java -cp .;test-classes org.testng.TestNG testng.xml` Note: Update the path_to_chrome_driver variable with your actual ChromeDriver executable path. The steps for the ChromeDriver executable and the command to run testng.xml can be removed, as they are not accurate. Instead, the step to run the testng.xml by right-clicking on it can be updated here. Overall, the AI agent generated the test code quickly, saving manual effort and providing a good starting point. However, review, validate, and refine the code before using it in a project. Best Practices While Using AI Agents The following best practices should be considered while using AI Agents to generate automation code: Improving prompt quality for reliable tests: Clear and detailed prompts help the AI generate stable test scripts. Mention requirements like framework names, naming conventions, and coding standards to improve structure and avoid errors.Adding assertions and validations: Always ensure the generated tests include proper assertions to validate expected outcomes. This helps confirm that the application behaves correctly instead of just performing actions.Review and refactor generated code before execution: Always review AI-generated code before running it. Refactoring helps remove redundancies, improve readability, and align the code with the project standards.Handling invalid/empty inputs: Validate inputs like test case files before processing them. Ensure you add proper steps and verification checks for empty or invalid data to prevent unexpected failures during execution.Integrating OpenAI API error handling: Proper error handling for API calls helps manage issues like timeouts, rate limits, or failed responses. Using try-catch blocks and logging error messages improves error readability. Exception handling helps interpret error messages in simple language and makes debugging easier.Logging and debugging tips: Adding logs at key steps helps track execution flow and quickly identify issues. Clear logging messages make debugging easier, especially when working with AI-generated outputs and API calls. Watch the step-by-step YouTube tutorial on How to Build an AI Agent to Generate Selenium WebDriver Tests in Java. Limitations and Challenges The following limitations and challenges can be faced while working with AI Agents for test script generation: AI-generated scripts may fail: AI-generated test scripts can be a starting point, but they may fail when the input prompt is unclear or missing important details. They can also break because of dynamic elements, incorrect UI assumptions, or environment-specific issues.Manual cleanup/selector accuracy: AI may generate locators that are not always reliable or optimal. Always review and refine locators in AI-generated scripts, preferring stable options like IDs or CSS selectors over complex XPaths.Test maintenance and flakiness considerations: AI-generated tests may be flaky if they don’t handle waits, dynamic content, or synchronization properly. Adding proper waits and improving test design can reduce flakiness. Regular maintenance is required to keep tests stable as the application evolves.Reviewing hallucinations: AI can sometimes generate incorrect or non-existent methods, elements, or logic. Review the code carefully to catch such hallucinations before execution. Validating against actual application behavior ensures the tests are accurate and usable. Conclusion To design an AI agent to generate test scripts, create a comprehensive prompt with clear rules, a structured output format, and strict guidelines to ensure consistent results. It’s also important to validate and refine the generated code and handle edge cases and failures gracefully. It all depends on the team’s decision; since AI is evolving, designing an AI agent that generates test automation scripts can also be a viable and flexible approach, especially when customization and control are important.
Enterprise applications commonly face multiple data challenges. Some data requires transactional integrity and relationships, while other data prioritizes fast, predictable access. Sessions, counters, rate limits, temporary state, often-accessed objects, and coordination data may not benefit from the complexity of a relational model. In these cases, a key-value database's simplicity becomes an architectural advantage. This simplicity is especially valuable in distributed and cloud-native systems, where latency, throughput, plus scalability directly shape user experience and infrastructure costs. A key-value database offers a focused approach: identify data by a key and retrieve or update it efficiently. The challenge is selecting a technology that delivers this performance while meeting the operational maturity, ecosystem support, and governance standards required for enterprise applications. Valkey meets these needs successfully. Originating from the Redis OSS lineage and developed as a vendor-neutral open-source project under the Linux Foundation, Valkey delivers a high-performance key-value platform suitable for caching, application state, messaging, and primary data storage. Beyond being another database option, it lets organizations explore how key-value persistence fits into modern enterprise architecture and lets Java applications use its benefits without tightly coupling to a specific datastore. Why Key-Value Databases Matter Key-value databases use a simple data model in which each unique key identifies a value. This simplicity is effective when applications can directly locate the required data. By enabling direct reads and writes, key-value databases typically deliver low latency, high throughput, and a horizontally scalable operational model. In enterprise systems, this model suits scenarios such as distributed sessions, caching, counters, rate limiting, feature flags, shopping carts, temporary workflow state, idempotency keys, leaderboards, and frequently accessed application data. These workloads prioritize fast access by identifier over joins, ad hoc queries, or complex relational constraints. The main architectural advantage of key-value databases is their specialization for specific access patterns, rather than universal speed or simplicity. When the primary requirement is to retrieve the current value for a given key, adding a more complex persistence model can introduce unnecessary overhead. As part of a polyglot persistence strategy, key-value stores enable architects to align the database model with the workload, rather than forcing all workloads into a single database. Putting Valkey Into Practice With Jakarta NoSQL A key advantage of using Valkey in enterprise Java is that it does not require a new programming model. With Jakarta NoSQL and Eclipse JNoSQL, Valkey serves as another key-value implementation behind a consistent API and mapping model. Domain annotations remain unchanged, so switching between key-value databases usually involves only updating the driver and its configuration, not rewriting the application. This abstraction is valuable architecturally. The application relies on the Jakarta NoSQL contract, while Eclipse JNoSQL manages integration with the database. Although database-specific features may introduce some coupling, applications that use the portable API can switch key-value implementations with minimal impact. For this article, we will use a simple Java SE example. This persistence layer can later support a REST API, messaging consumer, scheduled process, or other enterprise architecture without altering the core database interaction. Starting Valkey The first step is to make a Valkey instance available. Docker provides a convenient way to start one locally: Shell docker run --name valkey-instance \ -p 6379:6379 \ -d valkey/valkey:latest With Valkey running, add the Eclipse JNoSQL Valkey driver to the Jakarta NoSQL infrastructure, which includes CDI, Eclipse MicroProfile Config, and Jakarta JSON Processing. XML <dependency> <groupId>org.eclipse.jnosql.databases</groupId> <artifactId>jnosql-valkey</artifactId> <version>${jnosql.version}</version> </dependency> Configure the connection externally: Properties files jnosql.keyvalue.database=developers jnosql.valkey.port=6379 jnosql.valkey.host=localhost Since Eclipse JNoSQL integrates with Eclipse MicroProfile Config, you do not need to hard-code these values. They can be provided through configuration sources such as environment variables, in line with the Twelve-Factor App methodology. Mapping an Entity The mapping model for a key-value database is intentionally simple. Identify the class as an entity and specify the field that represents its key: Java @Entity public class User { @Id private String userName; private String name; private List<String> phones; // constructors, getters, setters... } Importantly, @Entity and @Id are part of the mapping abstraction, not Valkey itself. The domain model does not require Valkey-specific annotations. Using Jakarta NoSQL Eclipse JNoSQL provides KeyValueTemplate, a specialization of the Jakarta NoSQL Template API for key-value databases. This allows direct persistence and retrieval of entities: Java User user = User.builder() .phones(Arrays.asList("234", "432")) .username("username") .name("Name") .build(); KeyValueTemplate template = container.select(KeyValueTemplate.class).get(); User userSaved = template.put(user); System.out.println("User saved: " + userSaved); Optional<User> userFound = template.get("username", User.class); System.out.println("Entity found: " + userFound); For applications that prefer a repository abstraction, Eclipse JNoSQL integrates with Jakarta Data: Java @Repository public interface UserRepository extends CrudRepository<User, String> { } This approach allows the application code to focus more directly on domain operations: Java User user = User.builder() .phones(Arrays.asList("234", "432")) .username("username") .name("Name") .build(); UserRepository repository = container .select( UserRepository.class, DatabaseQualifier.ofKeyValue() ) .get(); repository.save(user); Optional<User> userFound = repository.findById("username"); System.out.println("User found: " + userFound); Notably, this code includes no Valkey-specific API in the entity or repository. Valkey is an infrastructure choice, while Jakarta NoSQL and Jakarta Data remain the application-facing abstractions. This separation guarantees the architecture remains reusable if the underlying key-value technology changes. Conclusion Key-value databases are highly effective for workloads that require direct access, low latency, and high throughput, rather than complex queries or relational navigation. This article examined how this model fits within enterprise architecture and how Valkey can integrate via Eclipse JNoSQL, allowing applications to avoid direct dependencies on vendor-specific APIs. By maintaining consistent entity mapping and using Jakarta NoSQL or Jakarta Data abstractions, switching key-value implementations becomes mainly a matter of infrastructure and configuration. This shift reflects a broader evolution in enterprise Java, as the platform expands its persistence capabilities beyond traditional relational databases. With Jakarta Persistence, Jakarta Data, Jakarta NoSQL, and tools like Eclipse JNoSQL, architects can choose the best data model for each workload while keeping familiar programming abstractions. Valkey enhances this ecosystem by providing a robust key-value option, making polyglot persistence both feasible and practical.
Oracle Database 23ai introduced the powerful DBMS_DEVELOPER package, giving developers and database administrators a streamlined way to access database object metadata in JSON format. This feature represents a significant advancement in how we interact with database schemas, offering a more structured and programmatic way to extract and analyze metadata compared to traditional dictionary views or the older DBMS_METADATA package. In this article, we'll explore the capabilities of DBMS_DEVELOPER, focusing on its GET_METADATA function through detailed examples and practical implementation scenarios. Understanding DBMS_DEVELOPER The DBMS_DEVELOPER package was designed specifically for modern application development patterns, where JSON has become a universal data exchange format. Rather than returning metadata as DDL statements (like DBMS_METADATA), this package returns structured JSON documents that can be easily parsed, processed, and integrated into applications or DevOps workflows. Key Benefits Structured data format: Returns metadata as JSON objects that can be easily parsed Programmatic access: Perfect for integration with applications and automation scripts Versioning capabilities: Built-in ETag mechanism for tracking object changesConfigurable detail levels: Ability to retrieve basic, typical, or comprehensive metadata Setting Up Our Environment Let's set up a sample schema to demonstrate the package functionality: SQL CREATE TABLE customers ( customer_id NUMBER(10) CONSTRAINT pk_customers PRIMARY KEY, first_name VARCHAR2(50) NOT NULL, last_name VARCHAR2(50) NOT NULL, email VARCHAR2(100) CONSTRAINT uk_customer_email UNIQUE, join_date DATE DEFAULT SYSDATE, status VARCHAR2(10) DEFAULT 'ACTIVE' ); CREATE INDEX idx_customer_name ON customers(last_name, first_name); CREATE OR REPLACE VIEW active_customers AS SELECT customer_id, first_name, last_name, email FROM customers WHERE status = 'ACTIVE'; GET_METADATA Basics The core function of the DBMS_DEVELOPER package is GET_METADATA, which returns metadata about database objects in JSON format. Let's start with a basic example: SQL -- Using JSON_SERIALIZE for formatted output SELECT JSON_SERIALIZE( DBMS_DEVELOPER.GET_METADATA(name => 'CUSTOMERS') PRETTY) AS metadata; The result is a structured JSON document containing comprehensive information about the table, including: Table name and schema Column definitions with data types and constraints Primary key, unique key, and foreign key information Index definitions An etag value representing the current state of the object This structured format makes it significantly easier to extract specific information programmatically compared to parsing DDL statements. NAME and SCHEMA Parameters The NAME and SCHEMA parameters work together to identify the specific database object. These parameters are case-sensitive and must match the object definition in the data dictionary. SQL -- Explicitly specifying schema SELECT JSON_SERIALIZE( DBMS_DEVELOPER.GET_METADATA( name => 'CUSTOMERS', schema => 'FINANCE') PRETTY) AS metadata; -- Using current schema (implicit) SELECT JSON_SERIALIZE( DBMS_DEVELOPER.GET_METADATA(name => 'CUSTOMERS') PRETTY) AS metadata; When the SCHEMA parameter is omitted, the function uses the current schema. This behavior provides flexibility when working with objects across different schemas in your database environment. OBJECT_TYPE Parameter The OBJECT_TYPE parameter allows you to explicitly specify the type of object you're retrieving metadata for. While often optional (as the database can infer the object type from the name), it becomes necessary in cases where name resolution alone is insufficient. Currently, `DBMS_DEVELOPER` supports three object types: TABLEINDEXVIEW Let's examine metadata for our index and view: SQL -- Retrieving index metadata SELECT JSON_SERIALIZE( DBMS_DEVELOPER.GET_METADATA( name => 'IDX_CUSTOMER_NAME', object_type => 'INDEX') PRETTY) AS metadata; -- Retrieving view metadata SELECT JSON_SERIALIZE( DBMS_DEVELOPER.GET_METADATA( name => 'ACTIVE_CUSTOMERS', object_type => 'VIEW') PRETTY) AS metadata; The OBJECT_TYPE parameter becomes particularly important when dealing with objects that share the same name but have different types, such as packages and package bodies. LEVEL Parameter The LEVEL parameter controls the amount of detail included in the JSON output. Oracle provides three levels: BASIC: Minimal informationTYPICAL: Standard level of detail (default)ALL: Comprehensive metadata This flexibility lets you balance concise output with detailed information based on your needs. SQL -- Basic level metadata SELECT JSON_SERIALIZE( DBMS_DEVELOPER.GET_METADATA( name => 'IDX_CUSTOMER_NAME', level => 'BASIC') PRETTY) AS metadata; -- All details SELECT JSON_SERIALIZE( DBMS_DEVELOPER.GET_METADATA( name => 'IDX_CUSTOMER_NAME', level => 'ALL') PRETTY) AS metadata; The output at the ALL level includes additional attributes such as segment information, compression settings, and physical storage details that aren't present at the BASIC level. ETAG Parameter One of the most powerful features of DBMS_DEVELOPER is the etag mechanism, which provides version tracking for database objects. The etag value changes whenever the object definition changes, making it invaluable for change detection. SQL -- Store the current etag value DECLARE v_metadata CLOB; v_etag VARCHAR2(100); BEGIN v_metadata := DBMS_DEVELOPER.GET_METADATA(name => 'ACTIVE_CUSTOMERS'); SELECT JSON_VALUE(v_metadata, '$.etag') INTO v_etag FROM dual; DBMS_OUTPUT.PUT_LINE('Current etag: ' || v_etag); END; / -- Modify the view CREATE OR REPLACE VIEW active_customers AS SELECT customer_id, first_name, last_name, email, join_date FROM customers WHERE status = 'ACTIVE'; -- Check if the object has changed using the stored etag SELECT JSON_SERIALIZE( DBMS_DEVELOPER.GET_METADATA( name => 'ACTIVE_CUSTOMERS', etag => 'A1B2C3D4E5F6G7H8I9J0') -- Previous etag value PRETTY) AS metadata; When you pass an ETag value that matches the current state of the object, the function returns an empty JSON document {}. If the object has changed, it returns the complete metadata with a new ETag value. Practical Scenario: Database Migration and Documentation Let's consider a practical scenario where DBMS_DEVELOPER proves invaluable: a large-scale database migration project with continuous schema changes. The Challenge You're leading a project to migrate a critical application database from on-premises to Oracle Cloud. The development team continues to make schema changes during the migration process, and you need to: Document the current state of all database objectsTrack changes between migration wavesValidate that objects were created correctly in the target environmentGenerate comprehensive documentation for compliance requirements The Solution Using DBMS_DEVELOPER, you can create a robust metadata management system: SQL CREATE TABLE schema_versions ( object_name VARCHAR2(128), object_type VARCHAR2(30), object_schema VARCHAR2(128), capture_date TIMESTAMP, etag VARCHAR2(100), metadata CLOB ); -- Procedure to capture all tables in a schema CREATE OR REPLACE PROCEDURE capture_schema_metadata(p_schema VARCHAR2) AS v_metadata CLOB; v_etag VARCHAR2(100); CURSOR c_objects IS SELECT object_name, object_type FROM all_objects WHERE owner = p_schema AND object_type IN ('TABLE', 'INDEX', 'VIEW'); BEGIN FOR obj IN c_objects LOOP BEGIN v_metadata := DBMS_DEVELOPER.GET_METADATA( name => obj.object_name, schema => p_schema, object_type => obj.object_type ); SELECT JSON_VALUE(v_metadata, '$.etag') INTO v_etag FROM dual; INSERT INTO schema_versions (object_name, object_type, object_schema, capture_date, etag, metadata) VALUES (obj.object_name, obj.object_type, p_schema, SYSTIMESTAMP, v_etag, v_metadata); COMMIT; DBMS_OUTPUT.PUT_LINE('Captured metadata for ' || obj.object_type || ' ' || p_schema || '.' || obj.object_name); EXCEPTION WHEN OTHERS THEN DBMS_OUTPUT.PUT_LINE('Error capturing ' || obj.object_type || ' ' || p_schema || '.' || obj.object_name || ': ' || SQLERRM); END; END LOOP; END; / This solution provides several key benefits: Efficient change tracking: Using etags to identify exactly which objects have changedStructured documentation: Storing metadata in JSON format for easy extraction of specific attributesHistorical record: Maintaining snapshots of schema evolution over timeValidation capabilities: Comparing source and target schemas during migration During migration, you can extend this system to compare environments: -- Procedure to compare object between environments CREATE OR REPLACE PROCEDURE compare_object( p_name VARCHAR2, p_type VARCHAR2, p_source_schema VARCHAR2, p_target_schema VARCHAR2, p_target_db VARCHAR2 ) AS v_source_metadata CLOB; v_target_metadata CLOB; v_source_etag VARCHAR2(100); v_target_etag VARCHAR2(100); BEGIN -- Get source metadata v_source_metadata := DBMS_DEVELOPER.GET_METADATA( name => p_name, schema => p_source_schema, object_type => p_type ); -- Get target metadata via database link EXECUTE IMMEDIATE 'SELECT DBMS_DEVELOPER.GET_METADATA( name => :1, schema => :2, object_type => :3 ) FROM dual@' || p_target_db INTO v_target_metadata USING p_name, p_target_schema, p_type; -- Extract etag values SELECT JSON_VALUE(v_source_metadata, '$.etag') INTO v_source_etag FROM dual; SELECT JSON_VALUE(v_target_metadata, '$.etag') INTO v_target_etag FROM dual; -- Compare and report IF v_source_etag = v_target_etag THEN DBMS_OUTPUT.PUT_LINE('Objects match exactly'); ELSE DBMS_OUTPUT.PUT_LINE('Objects differ - detailed comparison needed'); -- Further JSON comparison logic could be implemented here END; END; / Conclusion The DBMS_DEVELOPER package represents a significant advancement in Oracle's metadata management capabilities. By providing metadata in JSON format, Oracle has created a more developer-friendly interface that aligns with modern application architecture patterns. Key takeaways include: JSON-based metadata is more programmatically accessible than traditional DDL statements The etag mechanism provides a reliable way to track object changes Multiple detail levels allow you to retrieve just the information you need The package is particularly valuable for documentation, migration, and change tracking While currently limited to tables, indexes, and views, the DBMS_DEVELOPER package has tremendous potential for expansion in future Oracle releases. Database architects and developers should consider integrating this powerful tool into their workflows, particularly for projects involving schema documentation, migration, or programmatic metadata access. As databases continue to evolve toward more autonomous and programmable systems, tools like DBMS_DEVELOPER will become increasingly central to efficient database management practices.
After a decade of building and debugging large-scale data pipelines across financial services, payments processing, and analytics platforms, I can tell you that almost every slow Spark job I've investigated had the same root cause — and it wasn't the one the team thought it was. The default response when a Spark job is slow is to add more executor memory, increase the number of executors, or bump spark.sql.shuffle.partitions. Sometimes that helps. Usually it doesn't. What I've found, consistently, is that the real problems are structural — a join strategy mismatch that silently multiplies your intermediate dataset by ten times, a single slow task on a degraded node that holds an entire stage hostage, or a decrypt chain that re-reads source data six times when it only needed to read it once. This article is organized around five patterns I keep seeing across teams. Each one looks different on the surface but traces back to a misunderstanding of how Spark actually executes your code. For each pattern, I'll describe what it looks like, when it bites you, the failure mode, and how to fix it. Pattern 1: The OR Join That Quietly Multiplies Your Data What It Looks Like A join condition with an OR clause. Usually introduced when a business requirement adds a secondary matching rule — match on primary card number, or if the transaction is a virtual card transaction, match on the underlying physical PAN. The SQL looks reasonable. The engineer tests it on a sample, and it returns the right rows. When It Bites You At scale. With 100 million transaction rows and 50 million account rows, this query starts running for hours. The output size is also wrong — much larger than expected before DISTINCT trims it down. The Failure Mode Spark cannot use a hash join or sort-merge join when the join condition contains OR. It falls back to BroadcastNestedLoopJoin — for every row in the left table, scan every row in the right table. That's O(n x m). On real datasets, this produces an intermediate result in the hundreds of GB before any downstream filter runs. I've watched a pipeline that should produce 8 GB of output generate 400 GB of intermediate data because of exactly this pattern, taking a 20-minute job to 4 hours. You can verify this in 30 seconds: run df.explain(formatted) and look for BroadcastNestedLoopJoin in the physical plan. If you see it on a join involving any table over a few million rows, it's almost certainly unintentional. The Fix Split the join into two equi-join legs and UNION ALL the results: SQL -- Leg 1: primary match (equi-join — uses SortMergeJoin or BroadcastHashJoin) SELECT txn.*, acct.* FROM transactions txn JOIN accounts acct ON txn.card_number = acct.card_number UNION ALL -- Leg 2: fallback match, filtered scope only SELECT txn.*, acct.* FROM transactions txn JOIN accounts acct ON txn.fpan = acct.physical_pan WHERE txn.transaction_type = 'VIRTUAL' Each leg is a proper equi-join. Apply DISTINCT at the end to deduplicate rows that matched both. The performance difference is routinely an order of magnitude. Pattern 2: The Straggler Task That Nobody Notices Until It's Too Late What It Looks Like A stage that should take 10 minutes takes 3 hours. The Spark UI shows nearly all tasks completed quickly. One or two tasks are still running with a disproportionately long duration. When It Bites You Jobs running on shared YARN or cloud infrastructure where any node can have a bad disk, a noisy neighbor, or degraded network throughput. Also common in stages that call external services per partition — one slow API response can cause a single partition's tasks to take 100x longer than the others. The Failure Mode A stage doesn't complete until the last task completes. Not the median. Not p95. The absolute last one. If 2,200 tasks finish in under 2 minutes and one takes 3 hours and 7 minutes, the stage takes 3 hours and 7 minutes. The other 2,199 executors sit idle. This is the straggler problem, and it's distinct from data skew. The diagnostic: in the Stage detail view, check the task duration distribution. If MAX is dramatically higher than p99, that's a straggler (hardware or external service issue). If p75 is already much higher than p50, that's skew (data distribution issue). They require different fixes, and many teams treat them identically. The Fix For stragglers caused by degraded infrastructure, enable Spark speculation: Properties files spark.speculation=true spark.speculation.multiplier=3 # task must be 3x slower than median spark.speculation.quantile=0.9 # wait for 90% completion before speculating Speculation re-launches slow tasks on a different executor and uses whichever copy finishes first. The caveat: don't use this on stages that write to non-idempotent sinks. For read-heavy or compute-heavy stages — including external decryption calls — it's often the single most impactful config change you can make. Pattern 3: The df.rdd Decrypt Chain That Recomputes Everything Six Times What It Looks Like A pipeline that calls an external encryption or decryption service per record, implemented as a series of df.rdd.mapPartitions() calls, one per column that needs to be processed. When It Bites You When you have multiple columns to decrypt. Each .rdd call creates a new computation starting from the original DataFrame — Spark re-reads from source, re-executes all upstream joins and filters, and then runs the decryption for that column. With six columns to decrypt, you're doing that six times. The Failure Mode Two distinct sub-problems compound each other. First, going to RDD bypasses Catalyst entirely — no predicate pushdown, no column pruning, no Tungsten execution. Second, without a persist checkpoint before the chain, every decrypt call lineages all the way back to the source. I've seen this double the runtime of a job compared to the same pipeline with a single persist() before the decrypt chain. On top of that, the external call latency per partition is dominated by the number of HTTP round trips, not the payload size. Cutting your batch size in half doubles your request count and roughly doubles your wall-clock time for that stage. Most teams set an initial batch size and never revisit it. The Fix Two changes, applied together: Persist the input DataFrame before starting the decrypt chain. This means the join and filter logic runs once, and each decrypt call reads from the cached result.Increase the batch size for external calls. Test at several sizes — going from 20,000 to 40,000 records per batch often cuts stage time by 30-50% with no change to correctness. Scala val base = rawDf.filter(...).join(key1, ...).persist(StorageLevel.MEMORY_AND_DISK) val step1 = decryptColumn(base, secret1) // reads from cache val step2 = decryptColumn(step1, secret2) // reads from cache val step3 = decryptColumn(step2, secret3) // reads from cache Without persist, step2 re-executes everything step1 did from source. With persist, each step reads from the in-memory result of the previous. Pattern 4: The shuffle.partitions Setting That Nobody Updates What It Looks Like A job that works fine in staging — where data volumes are 10% of production — but runs slowly, spills to disk, or produces thousands of tiny output files in production. When It Bites You When the default spark.sql.shuffle.partitions=200 is left unchanged. 200 partitions made sense as a default for medium datasets but is almost always wrong at production scale — either too few (huge partitions, memory pressure) or too many (tiny partitions, scheduling overhead, small files problem). The Failure Mode Too few partitions means each executor handles a disproportionately large chunk of data. With 200 partitions on a 1 TB shuffle, each partition is 5 GB. That will spill to disk. Too many partitions means thousands of 1 MB tasks — the scheduling overhead becomes significant, and your output has thousands of tiny files that hurt downstream readers. With Adaptive Query Execution (AQE) enabled in Spark 3.2+, this problem largely manages itself. AQE merges small post-shuffle partitions automatically and can handle modest skew. But AQE can't help if it's disabled, and it can't fix the upstream causes of extreme skew. The Fix Enable AQE if you're on Spark 3.2+: Properties files spark.sql.adaptive.enabled=true spark.sql.adaptive.coalescePartitions.enabled=true spark.sql.adaptive.skewJoin.enabled=true If you need to set shuffle.partitions manually, target roughly 128-256 MB per partition post-shuffle. For a 500 GB shuffle, that means 2,000-4,000 partitions. Set it high and let AQE coalesce down — that's cheaper than setting it low and getting OOM errors. Pattern 5: The Incremental Job That Degrades Silently Over Time What It Looks Like A job that runs in 15 minutes when first deployed and runs in 4 hours six months later. No code changes. No obvious data quality issues. The team attributes it to data growth. When It Bites You When the job fails a few times in a row, and the recovery accumulates multiple windows' worth of data. Or when the watermark logic was designed for small windows but nobody anticipated that the underlying join tables would grow significantly. The Failure Mode Two separate causes, often confused. First, if the watermark is a single timestamp and the job has been failing, recovery runs can accumulate large backlogs. A job that normally processes 2 hours of data may need to process 48 hours on first successful recovery, with no change to the resource configuration. Second, growth in reference data (like an accounts table or lookup table used in a join) increases the size of every run regardless of whether the incremental input grew. I've seen a 30-minute job become a 3-hour job purely because the accounts table grew from 10 million rows to 80 million rows over 18 months, while the OR join condition (see Pattern 1) meant that growth was amplified into the intermediate result. The Fix Two design principles that pay off over the lifetime of the pipeline: Track processed partitions explicitly rather than using a single timestamp watermark. This makes recovery granular — you can replay specific missing partitions without re-processing everything after them.Add a fast-path no-op check before initializing the full Spark session. Check whether any new partitions exist first. A 5-second check that exits early is much better than a 2-minute executor startup that discovers there's nothing to process. For the reference table growth problem: if your lookup table grows significantly, revisit whether it can be broadcast (small enough to fit in executor memory) or whether the join itself needs to be redesigned. Quick Diagnostic Reference Use this table to map what you observe in the Spark UI to the likely pattern and first action to take: WHat you observeLikely patternconfirm withfirst action MAX task duration >> p99 Straggler (Pattern 2) Task timeline in Stage UI Enable spark.speculation p75 >> p50 task duration Data skew Input bytes per task Repartition on join key; AQE skewJoin BroadcastNestedLoopJoin in explain() OR join (Pattern 1) df.explain( formatted) Rewrite as UNION of equi-joins Stage runtime grows week on week; no code change Incremental accumulation or reference table growth (Pattern 5) Input bytes trend in History Server Audit watermark logic; check reference table size OOM errors or heavy disk spill Too few shuffle partitions (Pattern 4) Spill metrics in Stage UI Enable AQE or increase shuffle.partitions The Common Thread Every pattern here traces back to the same underlying issue: Spark is executing something different from what the engineer intended. The OR join was intended as a flexible matching rule; Spark turned it into a nested loop. The decrypt chain was intended as six independent transformations; Spark turned it into six full re-reads of source data. The incremental job was intended to process one window of data; without proper watermark design, it occasionally processes twelve. The Spark UI has everything you need to see this — task distribution, input and output sizes, physical plans, spill metrics. Most teams open it when something breaks and close it once they find the obvious error. Opening it proactively, forming a hypothesis, and then confirming or refuting it in the metrics is the practice that separates engineers who consistently improve pipeline performance from those who add executor memory and hope for the best. The mistake isn't choosing the wrong config. It's not understanding what Spark is actually doing with your code.
Data engineers managing batch SQL pipelines on Snowflake, BigQuery, and increasingly Databricks, and streaming pipelines on Apache Flink face a familiar problem: two toolchains, two skill sets, two CI/CD pipelines.dbt is now extending into stream processing. This post explains what that means in practice, why it matters for data engineering teams, and what a concrete implementation looks like with Apache Flink on Confluent Cloud. Data Streaming Meets the Lakehouse Data lakes promised to solve the enterprise data problem. The reality has been messier. Batch pipelines produce stale information, and analytical workloads run hours after the business event occurred. By the time a query runs, the window for action is often already closed. The lakehouse pattern has improved matters. Apache Iceberg has become the dominant open table format, supported across Snowflake, Databricks, BigQuery, and a growing number of query engines. Teams can run SQL analytics directly on data in object storage without duplicating it into a proprietary warehouse. But the lakehouse alone does not solve the real-time problem. Data still arrives as a batch, minutes or hours after the source event. That gap reflects a deeper architectural split. Data streaming with Apache Kafka and Flink is the operational layer: it handles critical SLAs, powers event-driven applications, and keeps business systems running in real time. The lakehouse is the analytical layer: it stores historical data for reporting, ML, and near real-time or batch analytics. These are two distinct workloads with different requirements regarding uptime, data loss, latency, and throughput. They need to coexist without forcing engineers to build and maintain two separate pipelines. How Kafka, Flink, and Iceberg Work Together That is what the combination of Apache Kafka, Apache Flink, and Apache Iceberg addresses. Kafka captures every event at the source and serves as the operational backbone for real-time systems. Flink processes and enriches data in motion, supporting both immediate operational decisions and the preparation of data for downstream analytics. Iceberg stores the result as a governed, queryable table for any analytical engine, whether that is Snowflake, BigQuery, or Databricks. A full treatment of this architecture, including schema evolution, compaction, and catalog integration, is covered here: Data Streaming Meets Lakehouse: Apache Iceberg for Unified Real-Time and Batch Analytics. The question is no longer whether streaming and lakehouse architectures can coexist. They already do. The question is how data engineering teams can work across both without maintaining separate toolchains. That is where dbt enters the picture. What Is dbt? dbt, the data build tool, is an open-source framework for SQL-based data transformation. A dbt model is a SQL SELECT statement saved as a file. dbt infers execution order from how models reference each other using ref(). The standard commands cover the full engineering workflow: dbt run executes the SQL against the target platform, dbt test validates data quality, and dbt docs generate produces a browsable documentation catalog. What made dbt successful is the discipline it brings to SQL work. Before dbt, transformation logic lived in scattered scripts and proprietary ETL tools. dbt replaced that with a code-first, version-controlled workflow with built-in lineage, testing, and documentation. Snowflake and BigQuery are where most dbt adoption lives today. Both are SQL-native and optimized for the ELT pattern dbt was built around. Redshift is a strong third platform in AWS environments. Databricks has seen growing dbt adoption more recently, driven by investments in serverless SQL Warehousing, but its roots are in Spark and Python, making it a newer entrant in the dbt ecosystem. dbt Labs crossed $100 million in ARR in early 2025, with over 5,000 paying customers. Around 90,000 dbt projects are running in production today. The Fivetran and dbt Labs merger, announced in October 2025, created a combined data infrastructure company with nearly $600 million in annual revenue — a clear signal that dbt has moved well beyond a popular open-source tool and into foundational enterprise data infrastructure. dbt Meets Apache Flink: One Workflow for Data Engineers Data engineering teams managing both batch and streaming today operate in two separate realities. Snowflake or BigQuery on one side: dbt models, version-controlled SQL, automated tests, generated docs. Apache Flink on the other: Terraform scripts, custom deployment code, or the Flink console. Skills and practices do not transfer between the two. That separation has a real cost. Streaming pipelines are harder to test, harder to document, and harder to hand over. Many teams compensate by keeping streaming logic minimal and pushing transformation work downstream into the warehouse, which reintroduces latency and undermines the point of streaming. The vision is straightforward: one SQL workflow for both. The engineer who builds dbt models on Snowflake or BigQuery should be able to apply the same approach to an Apache Flink streaming pipeline, without switching tools or rebuilding CI/CD from scratch. Two toolchains mean two testing strategies, two documentation systems, and two skill sets to hire and retain. Governance enforcement becomes inconsistent across the two environments. SQL is the shared foundation that makes this realistic. Flink SQL is mature and production-proven. Snowflake and BigQuery are SQL-native. Apache Iceberg tables are queryable via SQL across multiple engines. dbt wraps SQL with engineering discipline. The model files look the same. The ref() dependency resolution works the same way. Tests and documentation generation work through the same commands. Organizations do not need to hire separate Flink infrastructure specialists. The existing data engineering team can own both sides. Apache Iceberg connects the two worlds at the storage layer. A Flink pipeline writes structured, governed events into an Iceberg table in the organization's own S3 bucket. That same table is immediately readable by Snowflake, BigQuery, or Databricks without any additional ETL step. dbt can model data across the full pipeline: shaping it as it streams through Flink, and transforming it again when it lands in the warehouse for analytics. This is also a direct enabler of the Shift Left Architecture 2.0. The Shift Left approach moves data integration logic closer to the source, applying quality checks, enrichment, and governance in the streaming layer before data lands in the lakehouse. Until now, that required streaming-specific skills that most dbt-native teams did not have. dbt for Flink lowers that barrier considerably. The full architectural detail is covered here: The Shift Left Architecture 2.0: Operational, Analytical and AI Interfaces for Real-Time Data Products. Concrete Example: dbt on Confluent Cloud with Apache Flink The most concrete implementation available today is the dbt-confluent adapter, released by Confluent alongside the confluent-sql Python driver. Both are open source and available on PyPI and GitHub. Data engineers define streaming pipelines as dbt models and deploy them to Flink compute pools using the standard dbt run command. Getting started is a single step: pip install dbt-confluent Three materializations are supported: view for a virtual Flink SQL view over a Kafka topic, streaming_table for a continuous always-current result set, and streaming_source for defining a Kafka topic as a dbt source. Testing is deterministic, using Confluent Cloud's snapshot query capability to return bounded point-in-time results rather than silently passing on timeout. Documentation generation works through INFORMATION_SCHEMA integration, producing the same browsable catalog that Snowflake and BigQuery projects generate. The underlying confluent-sql driver is DB-API v2 compliant, meaning any compatible tool can connect directly to Confluent Cloud Flink: Airflow and Dagster for orchestration, Pandas for snapshot queries, Streamlit for live dashboards, and LangChain for AI agent workflows. For data engineers already working in dbt, this means the skills and practices built around Snowflake or BigQuery transfer directly to the streaming side of the architecture. The Data Engineer Owns Batch and Streaming with dbt The separation between batch and streaming engineering has always been more organizational than technical. Both worlds use SQL. Both require testing, documentation, and reliable deployment. The tools just never bridged the gap, so organizations staffed and operated two distinct engineering disciplines. dbt extending to Apache Flink changes that equation. The data engineer who runs dbt on Snowflake or BigQuery today can apply the same mental model, commands, and CI/CD pipeline to Flink streaming pipelines. No Flink infrastructure specialization required. They write SQL models, define tests, generate documentation, and deploy, exactly as they do for batch. The implication is straightforward. The investment in dbt skills and tooling now extends further into the architecture. Streaming can be adopted incrementally by the same data engineering teams already trusted for batch. One team, one tool, one governance standard, across both operational and analytical workloads. The Flink adapter for dbt is earlier in maturity compared to dbt on Snowflake or BigQuery, and teams should expect to work with an evolving ecosystem. But the foundation is solid, the direction is clear, and the core architectural components are already running in production at scale across multiple industries. The demand from data engineering teams is real and growing.
Founder,
Out of the Box Development, LLC