Source code for image_suggestions.wiki_indices
#!/usr/bin/env python3
# -*- coding: utf-8 -*-
"""Build full and delta datasets of boolean flags for all Wikipedias' search indices,
indicating whether an article has an image suggestion.
Flags follow weighted tags' syntax, namely ``recommendation.image/exists|1``
and ``recommendation.image_section/exists|1`` for :ref:`alis` and :ref:`slis` respectively.
The full dataset is stored in the :const:`image_suggestions.shared.SEARCH_INDEX_FULL_TABLE`,
and the delta in the :const:`image_suggestions.shared.SEARCH_INDEX_DELTA_TABLE`
`Hive <https://hive.apache.org/>`_ table of Wikimedia Foundation's
`Analytics Data Lake <https://wikitech.wikimedia.org/wiki/Analytics/Data_Lake>`_.
"""
import argparse
from pyspark.sql import DataFrame, SparkSession
from pyspark.sql import functions as F
from image_suggestions import queries, shared
MAIN_NAMESPACE_VALUE = 0
CONFIDENCE_COLUMN_NAME = 'confidence'
RECOMMENDATION_EXISTS_PREFIX = 'exists|'
# TODO consider one table for ALIS and one for SLIS
[docs]
def load_suggestions(
spark: SparkSession, hive_db: str, snapshot: str, target: str
) -> DataFrame: # pragma: no cover
"""Load image suggestions from :const:`image_suggestions.queries.SUGGESTIONS_TABLE`,
as output by :func:`image_suggestions.shared.save_suggestions`.
:param spark: an active Spark session
:param hive_db: a Hive database name
:param snapshot: a ``YYYY-MM-DD`` date
:param target: whether to load ALIS or SLIS. Accepted values: ``{'alis', 'slis'}``
:return: the dataframe of image suggestions
"""
return spark.sql(
queries.suggestions.format(
hive_db, snapshot, shared.get_section_index_constraint(target)
)
)
def parse_args() -> argparse.Namespace: # pragma: no cover
description = 'Build full and delta weighted tags for search indices'
parser = shared.build_base_arg_parser(description)
parser.add_argument(
'target', choices=['alis', 'slis'], help='Whether to build for ALIS or SLIS'
)
return parser.parse_args()
def main(args: argparse.Namespace, spark: SparkSession) -> None:
hive_db = args.hive_db
snapshot = args.snapshot
coalesce = args.coalesce
target = args.target
# Write full
suggestions = load_suggestions(spark, hive_db, snapshot, target)
exists_tags = build_exists_tags(target, suggestions)
shared.save_search_index_full(exists_tags, hive_db, snapshot, coalesce)
# Write delta (without Commons)
if target == 'alis':
target_tags = [shared.RECOMMENDATION_IMAGE_TAG]
elif target == 'slis':
target_tags = [shared.RECOMMENDATION_IMAGE_SECTION_TAG]
delta = shared.build_search_index_delta(spark, exists_tags, target_tags, snapshot)
shared.save_table(delta, hive_db, shared.SEARCH_INDEX_DELTA_TABLE, coalesce)
if __name__ == '__main__':
args = parse_args()
spark = shared.build_spark_session() # Pass 'dev' to switch environments
main(args, spark)
spark.stop()