""" Update status of feed in dynamo """ from collections import deque import datetime import boto from boto.dynamodb2.layer1 import DynamoDBConnection from boto.dynamodb2.table import Table from garcon.contrib.dynamo_feed_status import config STATUS_NOT_AVAILABLE = 'NOT_AVAILABLE' STATUS_DOWNLOADED = 'DOWNLOADED' STATUS_POPULATED_RAW_TABLE = 'POPULATED_RAW_TABLE' STATUS_INGESTED = 'INGESTED' STATUS_NOT_INGESTED = 'NOT_INGESTED' STATUS_INGESTED_TO_MYSQL = 'INGESTED_TO_MYSQL' STATUS_DEFAULT = None DYNAMODB_TABLE = config.feed_ingestion_table def get_dynamodb_connection(): """Return DynamoDB connection. Returns: DynamoDBConnection: connection instance. """ return DynamoDBConnection(host=config.host) def get_dynamodb_table(connection=None): """Return connection to table in DynamoDB. Args: connection (DynamoDBConnection): dynamo connection instance. Returns: boto.dynamodb2.table.Table: Table instance """ conn = connection or get_dynamodb_connection() return Table(DYNAMODB_TABLE, connection=conn) def get_attribute_name_for_file_status(file_name, date): """Get attribute name for file status Args: file_name (str): name of the file being extracted date (str): reporting date Returns: str: attribute name """ return '{}_status'.format( file_name.split('.')[0].replace('_{}'.format(date), '')) def set_status(feed_name, date, file_name, status=STATUS_NOT_AVAILABLE): """Set status of the ingestion job Args: feed_name (str): name of the feed date (str): reporting date of the ingestion job file_name (str): name of the file to be ingested status (str): status code constant """ attribute_name = get_attribute_name_for_file_status(file_name, date) # key(s) of the item to be updated key = { 'feed_name': {'S': feed_name}, 'date': {'S': date} } # we use a placeholder to avoid key validation error as the new key # being added is a filename and may have some forbidden chars in it placeholder = 'file_name_status' update_expression = 'SET #{} = :{}'.format(placeholder, placeholder) # defines what the placeholder should be replaced with using expression expression_attribute_names = { '#{}'.format(placeholder): '{}'.format(attribute_name) } # defines the value for that key expression_attribute_values = { ':{}'.format(placeholder): {'S': status} } get_dynamodb_connection().update_item( DYNAMODB_TABLE, key, update_expression=update_expression, expression_attribute_names=expression_attribute_names, expression_attribute_values=expression_attribute_values) def set_overall_status(feed_name, date, overall_status=None): """Set overall status Set overall status by looking at all the status for all files. If status is provided, set overall status with that. Args: feed_name (str): name of the feed date (str): reporting date of the ingestion job overall_status (str): overall_status code constant """ if not overall_status: overall_status = _determine_overall_status(feed_name, date) key = { 'feed_name': {'S': feed_name}, 'date': {'S': date} } update_expression = ( 'SET #status = :status') expression_attribute_names = { '#status': 'status' } expression_attribute_values = { ':status': {'S': overall_status} } get_dynamodb_connection().update_item( DYNAMODB_TABLE, key, update_expression=update_expression, expression_attribute_names=expression_attribute_names, expression_attribute_values=expression_attribute_values) def get_overall_status(feed_name, date): """Get status of the ingestion job Args: feed_name (str): Feed name date (str): Reporting date of the ingestion job Returns: mixed (False | status code): False if item not found. Status code if item is found. """ item = _get_item(feed_name, date) if item: return item['status'] return False def get_status(feed_name, date, file_name): """Get status of extracting the file Args: feed_name (str): Feed name date (str): Reporting date of the ingestion job file_name (str): name of the file that is being extracted Returns: str: status (NOT_AVAILABLE | DOWNLOADED) """ item = _get_item(feed_name, date) if item: attribute_name = get_attribute_name_for_file_status(file_name, date) return item[attribute_name] def _determine_overall_status(feed_name, date): """Set overall status of the feed Args: feed_name (str): name of the feed date (str): reporting date of the ingestion job Returns: str: overall status """ overall_status = STATUS_NOT_INGESTED item = _get_item(feed_name, date) assert item, 'No status found for feed {feed_name} and date {date}'. \ format(feed_name=feed_name, date=date) available_statuses = [STATUS_NOT_AVAILABLE, STATUS_DOWNLOADED] statuses = [s for k, s in item.items() if k not in [ 'feed_name', 'date', 'status']] num_of_files = len(statuses) for status in available_statuses: if num_of_files == statuses.count(status): overall_status = status return overall_status def _get_item(feed_name, date): """Get Dynamodb item Args: feed_name (str): Feed name date (str): Reporting date of the ingestion job Returns: mixed (False | status code): False if item not found. item if found. """ assert feed_name, 'feed_name is required.' assert date, 'date is required.' try: item = get_dynamodb_table().get_item( feed_name=feed_name, date=date) return item except boto.dynamodb2.exceptions.ItemNotFound: return False def delete_status(feed_name, date): """Delete item from dynamodb Args: feed_name (str): name fo the feed date (str): ingestion date """ get_dynamodb_table().delete_item(feed_name=feed_name, date=date) def get_missing_dates(feed_name, day_range=14, date_format='%Y-%m-%d'): """Find previous dates on which data was not ingested yet Args: feed_name (str): data feed name. day_range (int[optional]): number of days back from which search is being done. date_format (str[optional]): format of a date. Returns: deque: deque of dates string """ date_to_be_checked = datetime.date.today() - datetime.timedelta( days=day_range) dates = deque() while date_to_be_checked <= datetime.date.today(): status = get_overall_status( feed_name, date_to_be_checked.strftime(date_format)) if status != STATUS_INGESTED: dates.append(str(date_to_be_checked)) date_to_be_checked = date_to_be_checked + datetime.timedelta(days=1) return dates def find_non_processed_date(feed_name, available_dates): """Find files that have not been processed Args: feed_name (str): name of the feed available_dates (list): list of dates in str """ for available_date in available_dates: status = get_overall_status(feed_name, available_date) if status != STATUS_INGESTED: return available_date