diff --git a/luigi/contrib/batch.py b/luigi/contrib/batch.py index 4259c728be..d403eb1930 100644 --- a/luigi/contrib/batch.py +++ b/luigi/contrib/batch.py @@ -71,8 +71,10 @@ try: import boto3 + + _boto3_enabled = True except ImportError: - logger.warning("boto3 is not installed. BatchTasks require boto3") + _boto3_enabled = False class BatchJobException(Exception): @@ -88,6 +90,8 @@ def _random_id(): class BatchClient: def __init__(self, poll_time=POLL_TIME): + if not _boto3_enabled: + raise ImportError("boto3 is required for Batch functionality. Install it with: pip install boto3") self.poll_time = poll_time self._client = boto3.client("batch") self._log_client = boto3.client("logs") @@ -193,6 +197,11 @@ class BatchTask(luigi.Task): job_queue = luigi.OptionalParameter(default=None) poll_time = luigi.IntParameter(default=POLL_TIME) + def __init__(self, *args, **kwargs): + if not _boto3_enabled: + raise ImportError("boto3 is required for Batch functionality. Install it with: pip install boto3") + super().__init__(*args, **kwargs) + def run(self): bc = BatchClient(self.poll_time) job_id = bc.submit_job(self.job_definition, self.parameters, job_name=self.job_name, queue=self.job_queue) diff --git a/luigi/contrib/bigquery.py b/luigi/contrib/bigquery.py index a1b82aa01e..68ad84605a 100644 --- a/luigi/contrib/bigquery.py +++ b/luigi/contrib/bigquery.py @@ -32,8 +32,10 @@ try: import httplib2 from googleapiclient import discovery, errors, http + + _bigquery_enabled = True except ImportError: - logger.warning("BigQuery module imported, but google-api-python-client is not installed. Any BigQuery task will fail") + _bigquery_enabled = False else: RETRYABLE_ERRORS = (httplib2.HttpLib2Error, IOError, TimeoutError, BrokenPipeError) @@ -142,6 +144,8 @@ class BigQueryClient: """ def __init__(self, oauth_credentials=None, descriptor="", http_=None): + if not _bigquery_enabled: + raise ImportError("google-api-python-client is required for BigQuery functionality. Install it with: pip install google-api-python-client") # Save initialisation arguments in case we need to re-create client # due to connection timeout self.oauth_credentials = oauth_credentials @@ -398,6 +402,8 @@ def copy(self, source_table, dest_table, create_disposition=CreateDisposition.CR class BigQueryTarget(luigi.target.Target): def __init__(self, project_id, dataset_id, table_id, client=None, location=None): + if not _bigquery_enabled: + raise ImportError("google-api-python-client is required for BigQuery functionality. Install it with: pip install google-api-python-client") self.table = BQTable(project_id=project_id, dataset_id=dataset_id, table_id=table_id, location=location) self.client = client or BigQueryClient() diff --git a/luigi/contrib/bigquery_avro.py b/luigi/contrib/bigquery_avro.py index 4987942ed3..9a77078db6 100644 --- a/luigi/contrib/bigquery_avro.py +++ b/luigi/contrib/bigquery_avro.py @@ -11,8 +11,10 @@ try: import avro import avro.datafile + + _avro_enabled = True except ImportError: - logger.warning("bigquery_avro module imported, but avro is not installed. Any BigQueryLoadAvro task will fail to propagate schema documentation") + _avro_enabled = False class BigQueryLoadAvro(BigQueryLoadTask): @@ -30,6 +32,11 @@ class BigQueryLoadAvro(BigQueryLoadTask): source_format = SourceFormat.AVRO + def __init__(self, *args, **kwargs): + if not _avro_enabled: + raise ImportError("avro is required for BigQueryLoadAvro. Install it with: pip install avro-python3") + super().__init__(*args, **kwargs) + def _avro_uri(self, target): path_or_uri = target.uri if hasattr(target, "uri") else target.path return path_or_uri if path_or_uri.endswith(".avro") else path_or_uri.rstrip("/") + "/*.avro" diff --git a/luigi/contrib/datadog_metric.py b/luigi/contrib/datadog_metric.py index 4b6e9820ef..ff171f9f96 100644 --- a/luigi/contrib/datadog_metric.py +++ b/luigi/contrib/datadog_metric.py @@ -8,8 +8,10 @@ try: from datadog import api, initialize, statsd + + _datadog_enabled = True except ImportError: - logger.warning("Loading datadog module without datadog installed. Will crash at runtime if datadog functionality is used.") + _datadog_enabled = False class datadog(Config): @@ -24,6 +26,8 @@ class datadog(Config): class DatadogMetricsCollector(MetricsCollector): def __init__(self, *args, **kwargs): + if not _datadog_enabled: + raise ImportError("datadog is required for DatadogMetricsCollector. Install it with: pip install datadog") self._config = datadog(**kwargs) initialize(api_key=self._config.api_key, app_key=self._config.app_key, statsd_host=self._config.statsd_host, statsd_port=self._config.statsd_port) diff --git a/luigi/contrib/dataproc.py b/luigi/contrib/dataproc.py index f0a1ace682..0b74e9b8c0 100644 --- a/luigi/contrib/dataproc.py +++ b/luigi/contrib/dataproc.py @@ -19,11 +19,9 @@ DEFAULT_CREDENTIALS, _ = google.auth.default() authenticate_kwargs = gcp.get_authenticate_kwargs(DEFAULT_CREDENTIALS) _dataproc_client = discovery.build("dataproc", "v1", cache_discovery=False, **authenticate_kwargs) + _dataproc_enabled = True except ImportError: - logger.warning( - "Loading Dataproc module without the python packages googleapiclient & google-auth. \ - This will crash at runtime if Dataproc functionality is used." - ) + _dataproc_enabled = False def get_dataproc_client(): @@ -42,6 +40,14 @@ class _DataprocBaseTask(luigi.Task): dataproc_client = get_dataproc_client() + def __init__(self, *args, **kwargs): + if not _dataproc_enabled: + raise ImportError( + "google-api-python-client and google-auth are required for Dataproc functionality." + " Install them with: pip install google-api-python-client google-auth" + ) + super().__init__(*args, **kwargs) + class DataprocBaseTask(_DataprocBaseTask): """ diff --git a/luigi/contrib/docker_runner.py b/luigi/contrib/docker_runner.py index a3c9844e04..aee32fc498 100644 --- a/luigi/contrib/docker_runner.py +++ b/luigi/contrib/docker_runner.py @@ -48,9 +48,9 @@ import docker from docker.errors import APIError, ContainerError, ImageNotFound + _docker_enabled = True except ImportError: - logger.warning("docker is not installed. DockerTask requires docker.") - docker = None + _docker_enabled = False # TODO: may need to implement this logic for remote hosts # class dockerconfig(luigi.Config): @@ -143,6 +143,8 @@ def __init__(self, *args, **kwargs): - create a tmp dir - add the temp dir to the volume binds specified in the task """ + if not _docker_enabled: + raise ImportError("docker is required for DockerTask. Install it with: pip install docker") super(DockerTask, self).__init__(*args, **kwargs) self.__logger = logger diff --git a/luigi/contrib/dropbox.py b/luigi/contrib/dropbox.py index 5da75900b8..6aca7b7c18 100644 --- a/luigi/contrib/dropbox.py +++ b/luigi/contrib/dropbox.py @@ -33,10 +33,10 @@ import dropbox.dropbox_client import dropbox.exceptions import dropbox.files + + _dropbox_enabled = True except ImportError: - logger.warning( - "Loading Dropbox module without the python package dropbox (https://pypi.org/project/dropbox/). Will crash at runtime if Dropbox functionality is used." - ) + _dropbox_enabled = False def accept_trailing_slash_in_existing_dirpaths(func): @@ -74,6 +74,8 @@ def __init__(self, token, user_agent="Luigi", root_namespace_id=None): :param str token: Dropbox Oauth2 Token. See :class:`DropboxTarget` for more information about generating a token :param str root_namespace_id: Root namespace ID for interacting with Team Spaces """ + if not _dropbox_enabled: + raise ImportError("dropbox is required for Dropbox functionality. Install it with: pip install dropbox") if not token: raise ValueError("The token parameter must contain a valid Dropbox Oauth2 Token") @@ -292,6 +294,8 @@ def __init__(self, path, token, format=None, user_agent="Luigi", root_namespace_ """ + if not _dropbox_enabled: + raise ImportError("dropbox is required for Dropbox functionality. Install it with: pip install dropbox") super(DropboxTarget, self).__init__(path) if not token: diff --git a/luigi/contrib/ecs.py b/luigi/contrib/ecs.py index 3066e9a54b..c888abda1d 100644 --- a/luigi/contrib/ecs.py +++ b/luigi/contrib/ecs.py @@ -62,8 +62,9 @@ import boto3 client = boto3.client("ecs") + _boto3_enabled = True except ImportError: - logger.warning("boto3 is not installed. ECSTasks require boto3") + _boto3_enabled = False POLL_TIME = 2 @@ -138,6 +139,11 @@ class ECSTask(luigi.Task): task_def = luigi.OptionalParameter(default=None) cluster = luigi.Parameter(default="default") + def __init__(self, *args, **kwargs): + if not _boto3_enabled: + raise ImportError("boto3 is required for ECSTask. Install it with: pip install boto3") + super().__init__(*args, **kwargs) + @property def ecs_task_ids(self): """Expose the ECS task ID""" diff --git a/luigi/contrib/esindex.py b/luigi/contrib/esindex.py index 220a720171..61b1637b91 100644 --- a/luigi/contrib/esindex.py +++ b/luigi/contrib/esindex.py @@ -100,8 +100,9 @@ def docs(self): from elasticsearch.connection import Urllib3HttpConnection from elasticsearch.helpers import bulk + _elasticsearch_enabled = True except ImportError: - logger.warning("Loading esindex module without elasticsearch installed. Will crash at runtime if esindex functionality is used.") + _elasticsearch_enabled = False class ElasticsearchTarget(luigi.Target): @@ -129,6 +130,8 @@ def __init__(self, host, port, index, doc_type, update_id, marker_index_hist_siz :param extra_elasticsearch_args: extra args for Elasticsearch :type Extra: dict """ + if not _elasticsearch_enabled: + raise ImportError("elasticsearch is required for ElasticsearchTarget. Install it with: pip install elasticsearch") if extra_elasticsearch_args is None: extra_elasticsearch_args = {} diff --git a/luigi/contrib/ftp.py b/luigi/contrib/ftp.py index 07f745fa19..ec022ad70a 100644 --- a/luigi/contrib/ftp.py +++ b/luigi/contrib/ftp.py @@ -73,7 +73,7 @@ def _sftp_connect(self): try: import pysftp except ImportError: - logger.warning("Please install pysftp to use SFTP.") + raise ImportError("pysftp is required for SFTP functionality. Install it with: pip install pysftp") self.conn = pysftp.Connection(self.host, username=self.username, password=self.password, port=self.port, **self.pysftp_conn_kwargs) diff --git a/luigi/contrib/gcp.py b/luigi/contrib/gcp.py index c185f54aab..0c0bc64199 100644 --- a/luigi/contrib/gcp.py +++ b/luigi/contrib/gcp.py @@ -9,11 +9,10 @@ try: import google.auth import httplib2 + + _gcp_enabled = True except ImportError: - logger.warning( - "Loading GCP module without the python packages httplib2, google-auth. \ - This *could* crash at runtime if no other credentials are provided." - ) + _gcp_enabled = False def get_authenticate_kwargs(oauth_credentials=None, http_=None): @@ -25,6 +24,10 @@ def get_authenticate_kwargs(oauth_credentials=None, http_=None): Used by `gcs.GCSClient` and `bigquery.BigQueryClient` to initiate the API Client """ + if not _gcp_enabled: + raise ImportError( + "google-auth and google-auth-httplib2 are required for GCP functionality. Install them with: pip install google-auth google-auth-httplib2" + ) if oauth_credentials: authenticate_kwargs = {"credentials": oauth_credentials} elif http_: diff --git a/luigi/contrib/gcs.py b/luigi/contrib/gcs.py index 238c733f0b..08e283cf6f 100644 --- a/luigi/contrib/gcs.py +++ b/luigi/contrib/gcs.py @@ -40,11 +40,10 @@ try: import httplib2 from googleapiclient import discovery, errors, http + + _gcs_enabled = True except ImportError: - logger.warning( - "Loading GCS module without the python packages googleapiclient & google-auth. \ - This will crash at runtime if GCS functionality is used." - ) + _gcs_enabled = False else: RETRYABLE_ERRORS = (httplib2.HttpLib2Error, IOError) @@ -122,6 +121,8 @@ class GCSClient(luigi.target.FileSystem): """ def __init__(self, oauth_credentials=None, descriptor="", http_=None, chunksize=CHUNKSIZE, **discovery_build_kwargs): + if not _gcs_enabled: + raise ImportError("googleapiclient is required for GCS functionality. Install it with: pip install google-api-python-client") self.chunksize = chunksize authenticate_kwargs = gcp.get_authenticate_kwargs(oauth_credentials, http_) @@ -435,6 +436,8 @@ class GCSTarget(luigi.target.FileSystemTarget): fs = None def __init__(self, path, format=None, client=None): + if not _gcs_enabled: + raise ImportError("googleapiclient is required for GCS functionality. Install it with: pip install google-api-python-client") super(GCSTarget, self).__init__(path) if format is None: format = luigi.format.get_default_format() @@ -484,6 +487,8 @@ def __init__(self, path, format=None, client=None, flag="_SUCCESS"): :param flag: :type flag: str """ + if not _gcs_enabled: + raise ImportError("googleapiclient is required for GCS functionality. Install it with: pip install google-api-python-client") if format is None: format = luigi.format.get_default_format() diff --git a/luigi/contrib/kubernetes.py b/luigi/contrib/kubernetes.py index 0d68e45053..2f684a4ce0 100644 --- a/luigi/contrib/kubernetes.py +++ b/luigi/contrib/kubernetes.py @@ -46,8 +46,10 @@ from pykube.config import KubeConfig from pykube.http import HTTPClient from pykube.objects import Job, Pod + + _pykube_enabled = True except ImportError: - logger.warning("pykube is not installed. KubernetesJobTask requires pykube.") + _pykube_enabled = False class kubernetes(luigi.Config): @@ -62,6 +64,11 @@ class KubernetesJobTask(luigi.Task): __DEFAULT_POD_CREATION_INTERVAL = 5 _kubernetes_config = None # Needs to be loaded at runtime + def __init__(self, *args, **kwargs): + if not _pykube_enabled: + raise ImportError("pykube is required for KubernetesJobTask. Install it with: pip install pykube-ng") + super().__init__(*args, **kwargs) + def _init_kubernetes(self): self.__logger = logger self.__logger.debug("Kubernetes auth method: " + self.auth_method) diff --git a/luigi/contrib/mssqldb.py b/luigi/contrib/mssqldb.py index f89044e623..da6a6002ce 100644 --- a/luigi/contrib/mssqldb.py +++ b/luigi/contrib/mssqldb.py @@ -23,11 +23,10 @@ try: from pymssql import _mssql + + _pymssql_enabled = True except ImportError: - logger.warning( - "Loading MSSQL module without the python package pymssql. \ - This will crash at runtime if SQL Server functionality is used." - ) + _pymssql_enabled = False class MSSqlTarget(luigi.Target): @@ -56,6 +55,8 @@ def __init__(self, host, database, user, password, table, update_id): :param update_id: an identifier for this data set. :type update_id: str """ + if not _pymssql_enabled: + raise ImportError("pymssql is required for SQL Server functionality. Install it with: pip install pymssql") if ":" in host: self.host, self.port = host.split(":") self.port = int(self.port) diff --git a/luigi/contrib/mysqldb.py b/luigi/contrib/mysqldb.py index 0a1335193c..d5988d6f4a 100644 --- a/luigi/contrib/mysqldb.py +++ b/luigi/contrib/mysqldb.py @@ -25,11 +25,10 @@ try: import mysql.connector from mysql.connector import Error, errorcode + + _mysql_enabled = True except ImportError: - logger.warning( - "Loading MySQL module without the python package mysql-connector-python. \ - This will crash at runtime if MySQL functionality is used." - ) + _mysql_enabled = False class MySqlTarget(luigi.Target): @@ -56,6 +55,8 @@ def __init__(self, host, database, user, password, table, update_id, **cnx_kwarg :param cnx_kwargs: optional params for mysql connector constructor. See https://dev.mysql.com/doc/connector-python/en/connector-python-connectargs.html. """ + if not _mysql_enabled: + raise ImportError("mysql-connector-python is required for MySQL functionality. Install it with: pip install mysql-connector-python") if ":" in host: self.host, self.port = host.split(":") self.port = int(self.port) diff --git a/luigi/contrib/pai.py b/luigi/contrib/pai.py index 214b537d8d..d6bfa8e5a7 100644 --- a/luigi/contrib/pai.py +++ b/luigi/contrib/pai.py @@ -44,8 +44,9 @@ import requests as rs from requests.exceptions import HTTPError + _requests_enabled = True except ImportError: - logger.warning("requests is not installed. PaiTask requires requests.") + _requests_enabled = False def slot_to_dict(o): @@ -240,6 +241,8 @@ def __init__(self, *args, **kwargs): :param pai_url: The rest server url of PAI clusters, default is 'http://127.0.0.1:9186'. :param token: The token used to auth the rest server of PAI. """ + if not _requests_enabled: + raise ImportError("requests is required for PaiTask. Install it with: pip install requests") super(PaiTask, self).__init__(*args, **kwargs) self.__init_token() diff --git a/luigi/contrib/postgres.py b/luigi/contrib/postgres.py index 179ba5fb91..fd447565c3 100644 --- a/luigi/contrib/postgres.py +++ b/luigi/contrib/postgres.py @@ -68,10 +68,6 @@ def update_error_codes(): pass -if dbapi is None: - logger.warning("Loading postgres module without psycopg2 nor pg8000 installed. Will crash at runtime if postgres functionality is used.") - - def _is_pg8000_error(exception): try: return ( @@ -255,6 +251,8 @@ def connect(self): """ Get a DBAPI 2.0 connection object to the database where the table is. """ + if dbapi is None: + raise ImportError("psycopg2 or pg8000 is required for Postgres functionality. Install it with: pip install psycopg2") connection = dbapi.connect(host=self.host, port=self.port, database=self.database, user=self.user, password=self.password) connection.set_client_encoding("utf-8") return connection diff --git a/luigi/contrib/presto.py b/luigi/contrib/presto.py index 3b6bb1ccc5..4c10b4bb42 100644 --- a/luigi/contrib/presto.py +++ b/luigi/contrib/presto.py @@ -15,8 +15,10 @@ try: from pyhive.exc import DatabaseError from pyhive.presto import Connection, Cursor + + _pyhive_enabled = True except ImportError: - logger.warning("pyhive[presto] is not installed.") + _pyhive_enabled = False class presto(luigi.Config): # NOQA @@ -36,6 +38,8 @@ class PrestoClient: """ def __init__(self, connection, sleep_time=1): + if not _pyhive_enabled: + raise ImportError("pyhive is required for Presto functionality. Install it with: pip install pyhive[presto]") self.sleep_time = sleep_time self._connection = connection self._status = {"state": "initial"} diff --git a/luigi/contrib/redis_store.py b/luigi/contrib/redis_store.py index afe0717457..63ce15595b 100644 --- a/luigi/contrib/redis_store.py +++ b/luigi/contrib/redis_store.py @@ -26,8 +26,9 @@ try: import redis + _redis_enabled = True except ImportError: - logger.warning("Loading redis_store module without redis installed. Will crash at runtime if redis_store functionality is used.") + _redis_enabled = False class RedisTarget(Target): @@ -53,6 +54,8 @@ def __init__(self, host, port, db, update_id, password=None, socket_timeout=None :type expire: int """ + if not _redis_enabled: + raise ImportError("redis is required for RedisTarget. Install it with: pip install redis") self.host = host self.port = port self.db = db diff --git a/luigi/contrib/redshift.py b/luigi/contrib/redshift.py index db7a89d880..f2718995b0 100644 --- a/luigi/contrib/redshift.py +++ b/luigi/contrib/redshift.py @@ -32,7 +32,8 @@ import psycopg2 import psycopg2.errorcodes except ImportError: - logger.warning("Loading postgres module without psycopg2 installed. Will crash at runtime if postgres functionality is used.") + # Dependency check for psycopg2 is handled by contrib/postgres module + pass class _CredentialsMixin: diff --git a/luigi/contrib/s3.py b/luigi/contrib/s3.py index 87897089c7..d7b1687072 100644 --- a/luigi/contrib/s3.py +++ b/luigi/contrib/s3.py @@ -41,8 +41,10 @@ try: import botocore from boto3.s3.transfer import TransferConfig + + _boto3_enabled = True except ImportError: - logger.warning("Loading S3 module without the python package boto3. Will crash at runtime if S3 functionality is used.") + _boto3_enabled = False # two different ways of marking a directory # with a suffix in S3 @@ -72,6 +74,8 @@ class S3Client(FileSystem): DEFAULT_THREADS = 100 def __init__(self, aws_access_key_id=None, aws_secret_access_key=None, aws_session_token=None, **kwargs): + if not _boto3_enabled: + raise ImportError("boto3 is required for S3 functionality. Install it with: pip install boto3") options = self._get_s3_config() options.update(kwargs) if aws_access_key_id: diff --git a/luigi/contrib/salesforce.py b/luigi/contrib/salesforce.py index 27d81bcb73..a2bf087688 100644 --- a/luigi/contrib/salesforce.py +++ b/luigi/contrib/salesforce.py @@ -32,8 +32,10 @@ try: import requests + + _requests_enabled = True except ImportError: - logger.warning("This module requires the python package 'requests'.") + _requests_enabled = False def get_soql_fields(soql): @@ -235,6 +237,8 @@ class SalesforceAPI: API_NS = "{http://www.force.com/2009/06/asyncapi/dataload}" def __init__(self, username, password, security_token, sb_token=None, sandbox_name=None): + if not _requests_enabled: + raise ImportError("requests is required for Salesforce functionality. Install it with: pip install requests") self.username = username self.password = password self.security_token = security_token