mirror of
https://github.com/GoogleCloudPlatform/vertex-ai-samples.git
synced 2026-09-26 14:42:04 +00:00
refactor: update feature store vector search notebook use latest features. (#2909)
* Update feature store vector search notebook to use latest vertex SDK. * Run lint to fix format issue * Sleep for a few minutes to wait for DNS to be ready
This commit is contained in:
+14
-91
@@ -150,7 +150,6 @@
|
||||
"# Install the packages\n",
|
||||
"! pip3 install --upgrade --quiet google-cloud-aiplatform\\\n",
|
||||
" google-cloud-bigquery\\\n",
|
||||
" google-cloud-bigquery-datatransfer\\\n",
|
||||
" db-dtypes"
|
||||
]
|
||||
},
|
||||
@@ -528,15 +527,8 @@
|
||||
"outputs": [],
|
||||
"source": [
|
||||
"# First, create a dataset to keep the feature store source data if it does not already exist.\n",
|
||||
"BQ_DATASET_ID = (\n",
|
||||
" f\"featurestore_demo_{REGION.replace('-', '_')}\" # @param {type:\"string\"}\n",
|
||||
")\n",
|
||||
"create_bq_dataset(BQ_DATASET_ID, REGION)\n",
|
||||
"\n",
|
||||
"# The source data for this demo is located in the US region, so a temp dataset must also be located in the US region to transfer data to the desired region.\n",
|
||||
"TIMESTAMP = !date +%s\n",
|
||||
"BQ_TEMP_DATASET_ID = f\"fs_temp_{TIMESTAMP[0]}\"\n",
|
||||
"create_bq_dataset(BQ_TEMP_DATASET_ID, \"US\")"
|
||||
"BQ_DATASET_ID = \"featurestore_demo_us\" # @param {type:\"string\"}\n",
|
||||
"create_bq_dataset(BQ_DATASET_ID, \"US\")\n"
|
||||
]
|
||||
},
|
||||
{
|
||||
@@ -556,10 +548,11 @@
|
||||
},
|
||||
"outputs": [],
|
||||
"source": [
|
||||
"# 1. execute the query and keep results in temp dataset\n",
|
||||
"# Second, execute the query and store the results into a table\n",
|
||||
"BQ_TABLE_ID = \"publications_202304_small\" # @param {type:\"string\"}\n",
|
||||
"BQ_TEMP_TABLE_ID_FQN = f\"{PROJECT_ID}.{BQ_TEMP_DATASET_ID}.{BQ_TABLE_ID}\"\n",
|
||||
"job_config = bigquery.QueryJobConfig(destination=BQ_TEMP_TABLE_ID_FQN)\n",
|
||||
"BQ_TABLE_ID_FQN = f\"{PROJECT_ID}.{BQ_DATASET_ID}.{BQ_TABLE_ID}\"\n",
|
||||
"\n",
|
||||
"job_config = bigquery.QueryJobConfig(destination=BQ_TABLE_ID_FQN)\n",
|
||||
"query_job = bq_client.query(FEATURE_EXTRACT_QUERY_SMALL, job_config=job_config)\n",
|
||||
"\n",
|
||||
"try:\n",
|
||||
@@ -567,80 +560,9 @@
|
||||
"except Exception as e:\n",
|
||||
" # Table already exists\n",
|
||||
" print(\"Error: \", e.message)\n",
|
||||
"print(f\"Created table: {BQ_TEMP_TABLE_ID_FQN}\")"
|
||||
]
|
||||
},
|
||||
{
|
||||
"cell_type": "code",
|
||||
"execution_count": null,
|
||||
"metadata": {
|
||||
"id": "bab92f3c19e8"
|
||||
},
|
||||
"outputs": [],
|
||||
"source": [
|
||||
"import datetime\n",
|
||||
"import time\n",
|
||||
"\n",
|
||||
"from google.cloud import bigquery_datatransfer\n",
|
||||
"\n",
|
||||
"# 2. copy temp dataset to the desired region\n",
|
||||
"transfer_client = bigquery_datatransfer.DataTransferServiceClient()\n",
|
||||
"\n",
|
||||
"transfer_config = bigquery_datatransfer.TransferConfig(\n",
|
||||
" destination_dataset_id=BQ_DATASET_ID,\n",
|
||||
" display_name=f\"Copy BQ {BQ_TEMP_DATASET_ID} to {BQ_DATASET_ID}\",\n",
|
||||
" data_source_id=\"cross_region_copy\",\n",
|
||||
" params={\n",
|
||||
" \"source_project_id\": PROJECT_ID,\n",
|
||||
" \"source_dataset_id\": BQ_TEMP_DATASET_ID,\n",
|
||||
" },\n",
|
||||
" schedule=\"None\",\n",
|
||||
" schedule_options=bigquery_datatransfer.ScheduleOptions(disable_auto_scheduling=True),\n",
|
||||
")\n",
|
||||
"transfer_config = transfer_client.create_transfer_config(\n",
|
||||
" parent=transfer_client.common_project_path(PROJECT_ID),\n",
|
||||
" transfer_config=transfer_config,\n",
|
||||
")\n",
|
||||
"print(\n",
|
||||
" f\"Created data transfer config for dataset copy: {transfer_config.name}\"\n",
|
||||
")\n",
|
||||
"\n",
|
||||
"transfer_run = transfer_client.start_manual_transfer_runs(request=bigquery_datatransfer.StartManualTransferRunsRequest(\n",
|
||||
" parent=transfer_config.name,\n",
|
||||
" requested_run_time=datetime.datetime.now(tz=datetime.timezone.utc)\n",
|
||||
")).runs[0]\n",
|
||||
"\n",
|
||||
"print(\n",
|
||||
" f\"Started running dataset copy job: {transfer_run.name}\"\n",
|
||||
")\n",
|
||||
"\n",
|
||||
"job_completed = False\n",
|
||||
"while True:\n",
|
||||
" run = transfer_client.get_transfer_run(bigquery_datatransfer.GetTransferRunRequest(name=transfer_run.name))\n",
|
||||
" if run.state == bigquery_datatransfer.TransferState.SUCCEEDED:\n",
|
||||
" job_completed = True\n",
|
||||
" print(\"Dataset copy is complete.\")\n",
|
||||
" # clean up temp dataset.\n",
|
||||
" transfer_client.delete_transfer_config(\n",
|
||||
" request=bigquery_datatransfer.DeleteTransferConfigRequest(\n",
|
||||
" name=transfer_config.name\n",
|
||||
" )\n",
|
||||
" )\n",
|
||||
" bq_client.delete_dataset(\n",
|
||||
" bigquery.Dataset(f\"{PROJECT_ID}.{BQ_TEMP_DATASET_ID}\"),\n",
|
||||
" delete_contents=True,\n",
|
||||
" )\n",
|
||||
" # confirm table is created.\n",
|
||||
" BQ_TABLE_ID_FQN = f\"{PROJECT_ID}.{BQ_DATASET_ID}.{BQ_TABLE_ID}\"\n",
|
||||
" table = bq_client.get_table(bigquery.Table(BQ_TABLE_ID_FQN))\n",
|
||||
" DATA_SOURCE = f\"bq://{table}\"\n",
|
||||
" print(f\"Data source is: {DATA_SOURCE}\")\n",
|
||||
" elif run.state == bigquery_datatransfer.TransferState.FAILED:\n",
|
||||
" job_completed = True\n",
|
||||
" print(f\"Dataset copy is failed. \\n {run}\")\n",
|
||||
" if job_completed is True:\n",
|
||||
" break\n",
|
||||
" time.sleep(30)"
|
||||
"print(f\"Created table: {BQ_TABLE_ID}\")\n",
|
||||
"DATA_SOURCE = f\"bq://{BQ_TABLE_ID_FQN}\""
|
||||
]
|
||||
},
|
||||
{
|
||||
@@ -781,7 +703,7 @@
|
||||
"* A data source (BigQuery table or view URI or FeatureGroup/features ) synced to the `FeatureOnlineStore` instance for serving.\n",
|
||||
"* The cron schedule to run the sync pipeline.\n",
|
||||
"\n",
|
||||
"Within feature view creation, a sync job will be scheduled, either started immediately or following the cron schedule. In the sync job, data is exported to Cloud Bigtable, index is built and deployed to GKE cluster."
|
||||
"Within feature view creation, a sync job will be scheduled, either started immediately or following the cron schedule. In the sync job, data is exported, index is built and deployed to Feature Store backend."
|
||||
]
|
||||
},
|
||||
{
|
||||
@@ -794,7 +716,6 @@
|
||||
"source": [
|
||||
"FEATURE_VIEW_ID = \"feature_view_publications\" # @param {type: \"string\"}\n",
|
||||
"# A schedule will be created based on cron setting.\n",
|
||||
"# If cron is empty, an immediate schedule job will be started.\n",
|
||||
"CRON_SCHEDULE = \"TZ=America/Los_Angeles 00 13 11 8 *\" # @param {type: \"string\"}"
|
||||
]
|
||||
},
|
||||
@@ -831,10 +752,9 @@
|
||||
"\n",
|
||||
"index_config = utils.IndexConfig(\n",
|
||||
" embedding_column=EMBEDDING_COLUMN,\n",
|
||||
" dimensions=DIMENSIONS,\n",
|
||||
" crowding_column=CROWDING_COLUMN,\n",
|
||||
" filter_column=FILTER_COLUMNS,\n",
|
||||
" dimentions=DIMENSIONS,\n",
|
||||
" distance_measure_type=utils.DistanceMeasureType.DOT_PRODUCT_DISTANCE,\n",
|
||||
" filter_columns=FILTER_COLUMNS,\n",
|
||||
" algorithm_config=utils.TreeAhConfig(),\n",
|
||||
")\n",
|
||||
"\n",
|
||||
@@ -1057,6 +977,9 @@
|
||||
},
|
||||
"outputs": [],
|
||||
"source": [
|
||||
"# It will take some time for the DNS to be fully ready\n",
|
||||
"time.sleep(200)\n",
|
||||
"\n",
|
||||
"my_fv.search(\n",
|
||||
" entity_id=ENTITY_ID,\n",
|
||||
" neighbor_count=5,\n",
|
||||
|
||||
Reference in New Issue
Block a user