SIENTIAPDE-1110
Update Druid configuration in values.yaml to change port from 8081 to 8888. Refactor pydruid.py to use druid_engine for SQLAlchemy connections, enhancing clarity and consistency in data loading operations.
This commit is contained in:
@@ -19,7 +19,7 @@ class Druid(BaseActivity):
|
||||
|
||||
self.host = host
|
||||
self.port = port
|
||||
self.engine = create_engine(
|
||||
self.druid_engine = create_engine(
|
||||
f'druid://{self.host}:{self.port}/druid/v2/sql/')
|
||||
logger.info(
|
||||
f"Druid client initialized with host: {self.host}, port: {self.port}")
|
||||
@@ -51,10 +51,10 @@ class Druid(BaseActivity):
|
||||
f"Loading data from Druid: {datasource} with query: {query}"
|
||||
)
|
||||
|
||||
places = Table(datasource, MetaData(), autoload_with=self.engine)
|
||||
places = Table(datasource, MetaData(), autoload_with=self.druid_engine)
|
||||
stmt = select(places).where(text(query))
|
||||
|
||||
result = pd.read_sql(stmt, self.engine)
|
||||
result = pd.read_sql(stmt, self.druid_engine)
|
||||
|
||||
result["inserted_at"] = pd.to_datetime(result["__time"]).dt.strftime(
|
||||
"%Y-%m-%d %H:%M:%S.%f")
|
||||
|
||||
Reference in New Issue
Block a user