Hi, I tried to ingest metadata from PostgreSQL, bu...
# ingestion
b
Hi, I tried to ingest metadata from PostgreSQL, but I got the following error. Do you have any idea?
Copy code
$ ./scripts/datahub_docker.sh ingest -c ./postgres.yml
................................
OperationalError: (psycopg2.OperationalError) server didn't return client encoding'
The recipe (postgres.yml) is below.
Copy code
source: 
  type: postgres
  config: 
    host_port: <http://192.168.xxx.xxx:xxxx|192.168.xxx.xxx:xxxx>
    database: xxxx
    username: xxxx
    password: xxxx

sink: 
  type: "datahub-rest"
  config: 
    server: "<http://localhost:8080>"
m
Hi Ebu, Is that the only error you get or there is some stacktrace as well? If so, could you share it?
b
Since I am building datahub in an offline environment, I will write the log manually. The empty lines will be edited later.
Copy code
File "/usr/local/lib/python3.8/site-packages/datahub/entrypoints.py", line 95, in main
    sys.exit(datahub(standalone_mode=False, **kwargs))
File "/user/local/lib/python3.8/site-packages/click/core.py", line 1128, in __call__
    return self.main(*args, **kwargs)
File "/user/local/lib/python3.8/site-packages/click/core.py", line 1053, in main
    rv = self.invoke(ctx)
File "/user/local/lib/python3.8/site-packages/click/core.py", line 1659, in invoke
    return _process_result(sub_ctx.command.incoke(sub_ctx))
File "/user/local/lib/python3.8/site-packages/click/core.py", line 1659, in invoke
    return _process_result(sub_ctx.command.incoke(sub_ctx))
File "/user/local/lib/python3.8/site-packages/click/core.py", line 1395, in invoke
    return ctx.invoke(self.callback, **ctx.params)
File "/user/local/lib/python3.8/site-packages/click/core.py", line 754, in invoke
    return __callback(*args, **kqargs)    
File "/user/local/lib/python3.8/site-packages/datahub/cli/ingest_cli.py", line 74, in run
    pipeline.run()
File "/user/local/lib/python3.8/site-packages/datahub/ingestion/run/pipeline.py", line 148, in run
    for wu in itertools.islice(
File "/user/local/lib/python3.8/site-packages/datahub/ingestion/source/sql/sql_common.py", line 338, in get_workunits
    for inspector in self.get_inspectors():
File "/user/local/lib/python3.8/site-packages/datahub/ingestion/source/sql/sql_common.py", line 318, in get_inspectors
    with engine.connect() as conn:
File "/user/local/lib/python3.8/site-packages/sqlalchemy/engine/base.py", line 2263, in connect
    return self._connection_cls(self,**kwargs)
File "/user/local/lib/python3.8/site-packages/sqlalchemy/engine/base.py", line 104, in __init__
    else engine.raw_connection()
File "/user/local/lib/python3.8/site-packages/sqlalchemy/engine/base.py", line 2369, in raw_connection
    return self._wrap_pool_connect(
File "/user/local/lib/python3.8/site-packages/sqlalchemy/engine/base.py", line 2339, in _wrap_pool_connect
    Connection._handle_dbapi_exception_noconnection(
File "/user/local/lib/python3.8/site-packages/sqlalchemy/engine/base.py", line 1583, in _handle_dbapi_exception_noconnection
    util.raise_
File "/user/local/lib/python3.8/site-packages/sqlalchemy/util/compat.py", line 182, in raise_
        raise exception
File "/user/local/lib/python3.8/site-packages/sqlalchemy/engine/base.py", line 2336, in _wrap_pool_connect
    return fn()
File "/user/local/lib/python3.8/site-packages/sqlalchemy/pool/base.py", line 304, in unique_connection
    return _ConnectionFairy._checkout(self)
File "/user/local/lib/python3.8/site-packages/sqlalchemy/pool/base.py", line 778, in _checkout
    failry = _ConnectionRecord.chackout(pool)
File "/user/local/lib/python3.8/site-packages/sqlalchemy/pool/base.py", line 495, in checkout
    rec = pool._fo_get()
File "/user/local/lib/python3.8/site-packages/sqlalchemy/pool/impl.py", line 140, in _do_get
    self._dec_overflow()
File "/user/local/lib/python3.8/site-packages/sqlalchemy/util/langhelpers.py", line 68, in __exit__
    compat.raise_(
File "/user/local/lib/python3.8/site-packages/sqlalchemy/util/compat.py", line 182, in unique_connection
    raise exception
File "/user/local/lib/python3.8/site-packages/sqlalchemy/pool/impl.py", line 137, in _do_get
    return self._create_connection()
File "/user/local/lib/python3.8/site-packages/sqlalchemy/pool/base.py", line 309, in _create_connection
    return _ConnectionRecord(self)
File "/user/local/lib/python3.8/site-packages/sqlalchemy/pool/base.py", line 440, in __init__
    self.__connect(first_connect_check=True)
File "/user/local/lib/python3.8/site-packages/sqlalchemy/pool/base.py", line 661, in __connect
    pool.logger.debug("Error on connect(): %s", e)
File "/user/local/lib/python3.8/site-packages/sqlalchemy/util/langhelpers.py", line 68, in __exit__
    compat.raise_(
File "/user/local/lib/python3.8/site-packages/sqlalchemy/util/compat.py", line 182, in raise_
    raise exception
File "/user/local/lib/python3.8/site-packages/sqlalchemy/pool/base.py", line 656, in __connect
    connection = pool._invoke_creator(self)
File "/user/local/lib/python3.8/site-packages/sqlalchemy/engine/strategies.py", line 114, in connect
    return dialect.connect(*cargs, **cparams)
File "/user/local/lib/python3.8/site-packages/sqlalchemy/engine/default.py", line 508, in connect
    return self.dbapi.connect(*cargs, **cparams)
File "/user/local/lib/python3.8/site-packages/psycopg2/__init__.py", line 122, in connect
    conn = _connect (dsn, connection_factory=connection_factory, **kwasync)
@miniature-tiger-96062 Thank you for your patience. All logs are now available.
I’m trying to ingest from postgresql, but the goal is actually to ingest from Vertica. In this case, it seems that the plugin for vertica is not implemented yet, so do I need to use the source type for other sql alchemy databases?
I have to use a driver for vertica because psycopg2 gives me an error. Can you tell me the best driver for vertica and how to incorporate that driver into the datahub-ingestion container?
m
@breezy-controller-54597 are. you trying to ingest from Vertica using postgresql source ?
We have a generic source for sql databases that can connect using sql alchemy. You may want to check this out - https://datahubproject.io/docs/metadata-ingestion/source_docs/sqlalchemy
b
Yes. Because I thought that vertica is based on postgres. I will use Other SQLAlchemy source. Could you tell me how to incorporate the driver for vertica into the datahub-ingestion container?
I found “sqlalchemy-vertica-python” that use “v_catalog” instead of “pg_class”.
I installed “sqlalchemy-vertica-python” into “datahub-ingestion” container, but sqlalchemy use ‘sqlalchemy.dialects.postgresql’, not ‘vertica_python’.
m
Hi @breezy-controller-54597 , I don’t have a vetrica instance ready but did I did a quick search on how sql alchemy connects - engine = create_engine('vertica+pyodbc://username:password@mydsn')
So maybe try vertica+pyodbc?
b
Thank you for you suggestion. I tried ‘vertica+pyodbc’, but it is written by Python 2.7. I found ‘vertica+vertica_python’ that is forked from ‘vertica+pyodbc’.
It seems that sqlalchemy could load vertica+vertica_python, but it use the query for postgresql.
Copy code
File "/usr/local/lib/python3.8/site-packages/sqlalchemy/engine/case.py", line 1276, in _execute_context
    1186  def _execute_context(
    1187    self, dialect, constructor, statement, parameters, *args
    1188   ):
    (...)  
    1272        if fn(cursor, statement, parameters, context):
    1273          evt_handled = True
    1274          break
    1275      if not evt_handled:
--> 1276        self.dialect.do_execute(
    1277          cursor, statement, parameters, context
    ................................
    self = <sqlalchemy.engine.base.Connection object at 0x7fe8504dffd0>
    dialect = <sqlalchemy_vertica.dialect_vertica_python.VerticaDialect object at 0x7fe8504df820>
    constructor = <method 'DefaultExecutionContext._init_compiled' of <class 'sqlalchemy.dialects.postgresql.base.PGExecutionContext'> default.py:771>
    statement = "SELECT pg_get_view(c.oid) view_def FROM pg_class c JOIN pg_namespace n ON n.oid = c.relnamespace WHERE n.nspname = : schema AND c.relname = :view_name AND c.relkind IN ('v', 'm')"
    parameters = {'schema': 'public', 'view_name': 'V_FCT_ACCT_AGS'}
    args = (<sqlalchemy.dialects.postgresql.base.PGCompiler object at 0x7fe85041b1f0>, [{...}, ], )
    cursor = <vertica_python.vertica.cursor.Cursor object at 0x7fe8504875e0>
    context = <sqlalchemy.dialects.postgresql.base.PGExecutionContext object at 0x7fe85041beb0>
    evt_handled = False
    self.dialect.do_execute = <method 'DefaultDialect.do_execute' of <sqlalchemy_vertica.dialect_vertica_python.VerticaDialect object at 0x7fe8504df820> default.py:607>
Copy code
---- (full traceback above) ----
File "/usr/local/lib/python3.8/site-packages/datahub/entrypoint.py", line 95, in main
    sys.exit(datahub(standalone_mode=False, **kwargs))
File "/usr/local/lib/python3.8/site-packages/click/core.py", line 1128, in __call__
    return self.main(*args, **kwargs)
File "/usr/local/lib/python3.8/site-packages/click/core.py", line 1053, in main
    rv = self.invoke(ctx)
File "/usr/local/lib/python3.8/site-packages/click/core.py", line 1659, in invoke
    return _process_result(sub_ctx.command.invoke(sub_ctx))
File "/usr/local/lib/python3.8/site-packages/click/core.py", line 1659, in invoke
    return _process_result(sub_ctx.command.invoke(sub_ctx))
File "/usr/local/lib/python3.8/site-packages/click/core.py", line 1395, in invoke
    return ctx.invoke(self.callback, **ctx.params)
File "/usr/local/lib/python3.8/site-packages/click/core.py", line 754, in invoke
    return __callback(*args, **kwargs)
File "/usr/local/lib/python3.8/site-packages/datahub/cli/ingest_cli.py", line 74, in run
    pipeline.run()
File "/usr/local/lib/python3.8/site-packages/datahub/ingestion/run/pipeline.py", line 148, in run
    for wu in itertools.islice(
File "/usr/local/lib/python3.8/site-packages/datahub/ingestion/source/sql/sql_common.py", line 353, in get_workunits
    yield from self.loop_views(inspector, schema, sql_config)
File "/usr/local/lib/python3.8/site-packages/datahub/ingestion/source/sql/sql_common.py", line 573, in loop_views
    view_definition = inspector.get_view_definition(view, schema)
File "/usr/local/lib/python3.8/site-packages/sqlalchemy/engine/reflection.py", line 337, in get_view_definition
    return self.dialect.get_view_definition(
File "<string>", line 2, in get_view_definition
File "/usr/local/lib/python3.8/site-packages/sqlalchemy/engine/reflection.py", line 52, in cache
    ret = fn(self, con, *args, *kw)
File "/usr/local/lib/python3.8/site-packages/sqlalchemy/dialects/postgresql/base.py", line 3022, in get_view_definition
    view_def = connection.scalar(
File "/usr/local/lib/python3.8/site-packages/sqlalchemy/engine/base.py", line 941, in scalar
    return self.execute(object_, *multiparams, **params).scalar()
File "/usr/local/lib/python3.8/site-packages/sqlalchemy/engine/base.py", line 1011, in execute
    return meth(self, multiparams, params)
File "/usr/local/lib/python3.8/site-packages/sqlalchemy/sql/elements.py", line 298, in _execute_on_connection
    return connection._execute_clauseelement(self, multiparams, params)
File "/usr/local/lib/python3.8/site-packages/sqlalchemy/engine/base.py", line 1124, in _execute_clauseelement
    ret = self._execute_context(
File "/usr/local/lib/python3.8/site-packages/sqlalchemy/engine/base.py", line 1316, in _execute_context
    self._handle_dbapi_exception(
File "/usr/local/lib/python3.8/site-packages/sqlalchemy/engine/base.py", line 1510, in _handle_dbapi_exception
    util.raise_(
File "/usr/local/lib/python3.8/site-packages/sqlalchemy/util/compat.py", line 182, in raise_
    raise exception
File "/usr/local/lib/python3.8/site-packages/sqlalchemy/engine/base.py", line 1276, in _execute_context
    self.dialect.do_execute(
File "/usr/local/lib/python3.8/site-packages/sqlalchemy/engine/default.py", line 608, in do_execute
    cursor.execute(statement, parameters)
File "/usr/local/lib/python3.8/site-packages/vertica_python/vertica/cursor.py", line 222, in execute
    self._execute_simple_query(operation)
File "/usr/local/lib/python3.8/site-packages/vertica_python/vertica/cursor.py", line 606, in _execute_simple_query
    raise errors.QueryError.from_error_response(self._message, query)

ProgrammingError: (vertica_python.errors.MissingRelation) Severity: ERROR, Message: Relation "pg_class" does not exist, Sqlstate: 42V01, Routine: throwRelationDowsNotExist, File: /data/qb_workspaces/jenkins2/ReleaseBuilds/Hammerill/REL-10_0_1-x_hammermill/build/vertica/Catalog/CatalogLookup.cpp, Line: 3895, Error Code: 4566, SQL: "SELECT pg_get_viewdef(c.oid) view_def FROM pg_class c JOIN pg_namespace n ON n.oid - c.relnamespace WHERE n.nspname = 'xxxxx' AND c.relname = 'xxxxx' AND c.relkind IN ('v', 'm')"
[SQL: SELECT pg_get_viewdef(c.oid) view_def FROM pg_class c JOIN pg_namespace n ON n.oid - c.relnamespace WHERE n.nspname = 'xxxxx' AND c.relname = 'xxxxx' AND c.relkind IN ('v', 'm')]
[parameters: {'schema': 'xxxxx', 'view_name': 'xxxxx'}]
I think that it should be used below function in sqlalchemy-vertica-python for getting view information. https://github.com/bluelabsio/sqlalchemy-vertica-python/blob/99880f430ca4281eff845a47e424b33ff1520f12/sqla_vertica_python/vertica_python.py#L195 However, it seems that it is used below function in sqlalchemy/dialects/postgresql. https://github.com/sqlalchemy/sqlalchemy/blob/8d28bd8d4a0acc209884d29729020506a46f8bd4/lib/sqlalchemy/dialects/postgresql/base.py#L3584
m
I think the problem is that its trying to use postgresql dialect. What would make sense is we could try to just figure out how to connect to verica using sql_alchemy, outside of datahub
b
When called ‘get_inspectors’ in ‘/datahub/ingestion/source/sql/sql_common.py’, the inspector is set to ‘postgresql’. https://github.com/linkedin/datahub/blob/adf9d2ead771fd2661645b0b9e5260f739a7a023/[…]tadata-ingestion/src/datahub/ingestion/source/sql/sql_common.py It is assigned the inspector based on the url and config, but is it possible to specify the inspector explicitly in the recipe (yml)?
@miniature-tiger-96062 Should I submit it as an issue to github?