From 8c8d678241e03328fbf0e30f4106b19432565ff5 Mon Sep 17 00:00:00 2001 From: Mika Naylor Date: Fri, 11 Sep 2026 18:27:09 +0200 Subject: [PATCH] [FTAB-214] Update the lifecycle example to use new StatementHandle and statement lifecycle methods --- README.md | 29 +++-- examples/example_03_transforming_tables.py | 4 +- .../example_08_integration_and_deployment.py | 117 ++++++++++++------ pyproject.toml | 2 +- requirements.txt | 2 +- 5 files changed, 99 insertions(+), 55 deletions(-) diff --git a/README.md b/README.md index fcf969d..7abde8f 100644 --- a/README.md +++ b/README.md @@ -457,22 +457,29 @@ ConfluentTools.collect_materialized(table) ConfluentTools.print_materialized(table) ``` -### `ConfluentTools.get_statement_name` / `ConfluentTools.stop_statement` +### Managing statements with `StatementHandle` -Additional lifecycle methods are available to control statements on Confluent Cloud after they have -been submitted. +A `StatementHandle` controls a statement on Confluent Cloud after it has been submitted. Obtain one +from the `TableResult` of a statement you just submitted, or by name for a statement submitted +elsewhere (for example, from a separate CI/CD step). ```python -# On TableResult object -table_result = env.execute_sql("SELECT * FROM examples.marketplace.customers") -statement_name = ConfluentTools.get_statement_name(table_result) -ConfluentTools.stop_statement(table_result) +from confluent_pyflink.table.utils import StatementHandle -# Based on statement name -handle = ConfluentTools.get_statement_handle_by_name( - env, "table-api-2024-03-21-150457-36e0dbb2e366-sql" -) +# From the TableResult of a submitted statement +table_result = env.execute_sql("SELECT * FROM examples.marketplace.customers") +handle = StatementHandle.from_table_result(table_result) +print(handle.get_name()) handle.stop() + +# Or look up an existing statement by name +handle = StatementHandle.from_name(env, "table-api-2024-03-21-150457-36e0dbb2e366-sql") +handle.resume() +handle.delete() + +# Inspect non-fatal warnings raised while processing the statement +for warning in handle.get_warnings(): + print(warning.severity.value, warning.reason, warning.message) ``` ### Confluent Table Descriptor diff --git a/examples/example_03_transforming_tables.py b/examples/example_03_transforming_tables.py index 2e65ab7..7893186 100644 --- a/examples/example_03_transforming_tables.py +++ b/examples/example_03_transforming_tables.py @@ -18,7 +18,7 @@ from confluent_pyflink.table import TableEnvironment, DataTypes from confluent_pyflink.table.utils import ConfluentSettings -from confluent_pyflink.table.expressions import col, row +from confluent_pyflink.table.expressions import col, row, with_all_columns # A table program example that demos how to transform data with the Table object. @@ -34,7 +34,7 @@ def run(): # pipeline. No execution happens until execute() is called! # Read from tables like 'orders' - orders = env.from_path("orders") + orders = env.from_path("orders").select(with_all_columns()) # Or mock tables with values customers = env.from_elements( diff --git a/examples/example_08_integration_and_deployment.py b/examples/example_08_integration_and_deployment.py index 66673cd..e707ac1 100644 --- a/examples/example_08_integration_and_deployment.py +++ b/examples/example_08_integration_and_deployment.py @@ -16,11 +16,11 @@ # limitations under the License. ################################################################################ -import sys +import argparse import uuid from confluent_pyflink.table import TableEnvironment -from confluent_pyflink.table.utils import ConfluentSettings, ConfluentTools -from confluent_pyflink.table.expressions import lit +from confluent_pyflink.table.utils import ConfluentSettings, ConfluentTools, StatementHandle +from confluent_pyflink.table.expressions import lit, with_all_columns # NOTE: This example requires write access to a Kafka cluster. Fill out the # given variables below with target catalog/database if this is fine for you. @@ -38,8 +38,8 @@ # The following SQL will be tested on a finite subset of data before # it gets deployed to production. # In production, it will run on unbounded input. -# The '%s' parameterizes the SQL for testing. -SQL = "SELECT brand, COUNT(*) AS vendors FROM ProductsMock %s GROUP BY brand" +# The '{hints}' field parameterizes the SQL for use during testing. +SQL = "SELECT brand, COUNT(*) AS vendors FROM ProductsMock {hints} GROUP BY brand" # An example that illustrates how to embed a table program into a CI/CD @@ -57,8 +57,12 @@ # python example_08_integration_and_deployment test # python example_08_integration_and_deployment deploy # -# NOTE: The example submits an unbounded background statement. Make sure -# to stop the statement in the Web UI afterward to clean up resources. +# NOTE: The deploy phase submits an unbounded background statement. Clean it up +# afterward with the operate modes, which manage a statement by name: +# +# python example_08_integration_and_deployment stop +# python example_08_integration_and_deployment resume +# python example_08_integration_and_deployment delete # # The complete CI/CD workflow performs the following steps: # - Create Kafka table 'ProductsMock' and 'VendorsPerBrand'. @@ -69,30 +73,33 @@ # - Deploy an unbounded version of the tested SQL that writes into # 'VendorsPerBrand'. def run(args=None): - """Process command line arguments.""" - if not args: - args = sys.argv[1:] - - if len(args) == 0: - print("No mode specified. Possible values are 'setup', 'test', or 'deploy'.") - exit(1) + parser = argparse.ArgumentParser( + prog="example_08_integration_and_deployment", + description="CI/CD integration and deployment example.", + ) + sub = parser.add_subparsers(dest="mode", required=True) + sub.add_parser("setup", help="create tables and fill the source with data") + sub.add_parser("test", help="run the SQL on bounded data and check the result") + sub.add_parser("deploy", help="submit the unbounded statement") + for action in ("stop", "resume", "delete"): + p = sub.add_parser(action, help=f"{action} a deployed statement by name") + p.add_argument("statement_name", help=f"name of the statement to {action}") - mode = args[0] + parsed_args = parser.parse_args(args) settings = ConfluentSettings() env = TableEnvironment.create(settings) env.use_catalog(TARGET_CATALOG) env.use_database(TARGET_DATABASE) - if mode == "setup": + if parsed_args.mode == "setup": _set_up_program(env) - elif mode == "test": + elif parsed_args.mode == "test": _test_program(env) - elif mode == "deploy": + elif parsed_args.mode == "deploy": _deploy_program(env) else: - print("Unknown mode: " + mode) - exit(1) + _manage_statement(env, parsed_args.mode, parsed_args.statement_name) # -------------------------------------------------------------------------- @@ -101,23 +108,28 @@ def run(args=None): def _set_up_program(env: TableEnvironment): print("Running setup...") - print("Creating table..." + SOURCE_TABLE) + print(f"Creating table {SOURCE_TABLE}...") # Create a mock table that has exactly the same schema as the example # `products` table. # The LIKE clause is very convenient for this task which is why we use SQL # here. Since we use little data, a bucket of 1 is important to satisfy the - # `scan.bounded.mode` during testing. - env.execute_sql( - "CREATE TABLE IF NOT EXISTS `%s`\n" - "DISTRIBUTED INTO 1 BUCKETS\n" - "LIKE `examples`.`marketplace`.`products` (EXCLUDING OPTIONS)" % SOURCE_TABLE - ) + # `scan.bounded.mode` during testing. read-uncommitted makes the freshly filled + # rows visible immediately, rather than waiting for the exactly-once checkpoint + # commit. + env.execute_sql(f""" + CREATE TABLE IF NOT EXISTS `{SOURCE_TABLE}` + DISTRIBUTED INTO 1 BUCKETS + WITH ('kafka.consumer.isolation-level' = 'read-uncommitted') + LIKE `examples`.`marketplace`.`products` (EXCLUDING OPTIONS) + """) print("Start filling table...") # Let Flink copy generated data into the mock table. Note that the # statement is unbounded and submitted as a background statement by default. - pipeline_result = env.from_path("`examples`.`marketplace`.`products`").execute_insert( - SOURCE_TABLE + pipeline_result = ( + env.from_path("`examples`.`marketplace`.`products`") + .select(with_all_columns()) + .execute_insert(SOURCE_TABLE) ) print("Waiting for at least 200 elements in table...") @@ -136,13 +148,13 @@ def _set_up_program(env: TableEnvironment): # still needs a manual stop. ConfluentTools.stop_statement(pipeline_result) - print("Creating table..." + TARGET_TABLE) + print(f"Creating table {TARGET_TABLE}...") # Create a table for storing the results after deployment. - env.execute_sql( - "CREATE TABLE IF NOT EXISTS `%s` \n" - "(brand STRING, vendors BIGINT, PRIMARY KEY(brand) NOT ENFORCED)\n" - "DISTRIBUTED INTO 1 BUCKETS" % TARGET_TABLE - ) + env.execute_sql(f""" + CREATE TABLE IF NOT EXISTS `{TARGET_TABLE}` + (brand STRING, vendors BIGINT, PRIMARY KEY(brand) NOT ENFORCED) + DISTRIBUTED INTO 1 BUCKETS + """) # ----------------------------------------------------------------------------- @@ -164,7 +176,7 @@ def _test_program(env: TableEnvironment): ) print("Requesting test data...") - result = env.execute_sql(SQL % dynamicOptions) + result = env.execute_sql(SQL.format(hints=dynamicOptions)) rows = ConfluentTools.collect_materialized(result) print("Test data:") @@ -190,18 +202,43 @@ def _deploy_program(env: TableEnvironment): # It is possible to give a better statement name for deployment but make sure # that the name is unique within environment and region. - statement_name = "vendors-per-brand-" + str(uuid.uuid4()) + statement_name = f"vendors-per-brand-{uuid.uuid4()}" ConfluentTools.set_statement_name(env, statement_name) # Execute the SQL without dynamic options. # The result is unbounded and piped into the target table. - result = env.sql_query(SQL % "").execute_insert(TARGET_TABLE) + result = env.sql_query(SQL.format(hints="")).execute_insert(TARGET_TABLE) + + # A handle manages the submitted statement on Confluent Cloud. + handle = StatementHandle.from_table_result(result) # The API might add suffixes to manual statement names such as '-sql' or # '-api'. For the final submitted name, use the provided tools. - finalName = ConfluentTools.get_statement_name(result) + print(f"Statement has been deployed as: {handle.get_name()}") + + # Warnings surface non-fatal issues (such as deprecations) that did not stop the + # statement from being submitted. + warnings = handle.get_warnings() + for warning in warnings: + print(f" warning [{warning.severity.value}] {warning.reason}: {warning.message}") + + +# ---------------------------------------------------------------------------- +# Operate Phase +# ---------------------------------------------------------------------------- +def _manage_statement(env: TableEnvironment, action: str, statement_name: str): + print(f"Running {action}...") + + handle = StatementHandle.from_name(env, statement_name) + + if action == "stop": + handle.stop() + elif action == "resume": + handle.resume() + elif action == "delete": + handle.delete() - print("Statement has been deployed as: " + finalName) + print(f"{action}: {handle.get_name()}") if __name__ == "__main__": diff --git a/pyproject.toml b/pyproject.toml index 81caf71..ad6474b 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -9,7 +9,7 @@ license = {text = "Apache-2.0"} readme = "README.md" requires-python = ">=3.9,<3.12" dependencies = [ - "confluent-pyflink>=2.3.2", + "confluent-pyflink>=2.3.3", ] [dependency-groups] diff --git a/requirements.txt b/requirements.txt index 28f0be4..e7e0969 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1 +1 @@ -confluent-pyflink==2.3.2 +confluent-pyflink==2.3.3