"""Calls the social platform collectors.""" from datetime import datetime from multiprocessing.pool import ThreadPool from social_analytics.connectors import loggly from social_analytics.connectors import sentry from social_analytics.logic import collector from social_analytics.models import artist_url from social_analytics.models import social_profile logger = loggly.get_current_logger() def collect_facts(): """Run collectors for social platform metrics.""" # Get all profiles that are eligible for collection logger.info( ('Data collection process started at {0}'.format(str(datetime.now())))) profiles = (artist_url.get_social_profiles_for_collection()) if not profiles: return profiles = profiles.message['items'] logger.info( ('About to start collection for {0} profiles'.format(len(profiles)))) tokens = collector.get_tokens() # Create object for each social profile that contains tokens collection_params = [ { 'social_profile': profile, 'tokens': tokens } for profile in profiles ] with ThreadPool(40) as pool: processed_profiles = [] errors = [] errors_upload_raw_data = None errors_upload_metrics_data = None updated_items = None try: # Fetch metrics for each profile processed_profiles = pool.map(collector.collect, collection_params) logger.info( ('{0} Profiles have been collected and processed.'.format( len(processed_profiles)))) # Save the raw json to s3 logger.info('About to upload the raw data to s3') errors_upload_raw_data = collector.upload_raw_json_to_s3( processed_profiles) logger.info('Finished uploading raw data to s3') # Save the metrics json to s3 logger.info('About to upload processed data to s3') errors_upload_metrics_data = collector.upload_metrics_json_to_s3( processed_profiles) logger.info('Finished uploading processed data to s3') errors = list( filter(lambda x: 'error' in x.keys(), processed_profiles)) # Update the last_collected_date value for all collected profiles. collected_ids = [ x['social_profile_id'] for x in processed_profiles if x not in errors] logger.info(('The collected ids are: {0}'.format(collected_ids))) updated_items = ( social_profile.update_social_profiles_last_collected_date( collected_ids)) except Exception as exc: errors.append(str(exc)) logger.error( ('Exception occurred while collecting data {0}'.format( str(exc)))) sentry_client = sentry.get_client() sentry_client.captureMessage( message='Errors during collection', stack=True, extra={ 'message': 'Collection errors', 'errors': errors, 'status': 204}) results = { 'expected_profiles_count': len(processed_profiles), 'actual_profiles_count': (len(processed_profiles) - len(errors)), 'profiles_data': processed_profiles, 'errors_count': len(errors), 'errors_upload_raw_data': errors_upload_raw_data, 'errors_upload_metrics_data': errors_upload_metrics_data, 'errors_logging': errors, 'expected_updated_profiles': len(processed_profiles), 'actual_updated_profiles': updated_items } logger.info('The data collection finished at {}'.format(str(datetime))) logger.info('The results of the collection are: {0}'.format(results)) return results