Skip to content

Commit

Permalink
Add support for external engines
Browse files Browse the repository at this point in the history
  • Loading branch information
ryannedolan committed Jan 16, 2025
1 parent aae0587 commit be410c9
Show file tree
Hide file tree
Showing 17 changed files with 1,169 additions and 34 deletions.
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -83,7 +83,7 @@ Once the Flink deployment pod has STATUS 'Running', you can forward port 8081 an
to access the Flink dashboard.

```
$ kubectl port-forward basic-session-deployment-7b94b98b6b-d6jt5 8081 &
$ kubectl port-forward svc/basic-session-deployment-rest 8081 &
```

See the [Flink SQL Gateway Documentation](https://nightlies.apache.org/flink/flink-docs-release-1.18/docs/dev/table/sql-gateway/overview/)
Expand Down
8 changes: 8 additions & 0 deletions deploy/samples/flinkengine.yaml
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
apiVersion: hoptimator.linkedin.com/v1alpha1
kind: Engine
metadata:
name: flink-engine
spec:
url: jdbc:flink://localhost:8083
dialect: Flink

1 change: 1 addition & 0 deletions gradle/libs.versions.toml
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ flink-clients = "org.apache.flink:flink-clients:1.18.1"
flink-connector-base = "org.apache.flink:flink-connector-base:1.18.1"
flink-core = "org.apache.flink:flink-core:1.18.1"
flink-csv = "org.apache.flink:flink-csv:1.18.1"
flink-jdbc = "org.apache.flink:flink-sql-jdbc-driver-bundle:1.18.1"
flink-streaming-java = "org.apache.flink:flink-streaming-java:1.18.1"
flink-table-api-java = "org.apache.flink:flink-table-api-java:1.18.1"
flink-table-api-java-bridge = "org.apache.flink:flink-table-api-java-bridge:1.18.1"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -10,4 +10,6 @@ public interface Engine {
DataSource dataSource();

SqlDialect dialect();

String url();
}
1 change: 1 addition & 0 deletions hoptimator-cli/build.gradle
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@ dependencies {
implementation libs.calcite.core
implementation libs.sqlline
implementation libs.slf4j.simple
implementation libs.flink.jdbc
}

publishing {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,8 @@
import com.linkedin.hoptimator.Engine;
import com.linkedin.hoptimator.SqlDialect;

import java.util.Objects;

import org.apache.calcite.adapter.jdbc.JdbcSchema;


Expand All @@ -17,7 +19,7 @@ public class K8sEngine implements Engine {

public K8sEngine(String name, String url, SqlDialect dialect, String driver) {
this.name = name;
this.url = url;
this.url = Objects.requireNonNull(url, "url");
this.dialect = dialect;
this.driver = driver;
}
Expand All @@ -33,8 +35,13 @@ public DataSource dataSource() {
return JdbcSchema.dataSource(url, driver, null, null);
}

@Override
public String url() {
return url;
}

@Override
public SqlDialect dialect() {
return SqlDialect.FLINK; // TODO fix hardcoded dialect
return dialect;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@
import java.util.Locale;
import java.util.Optional;
import java.util.stream.Collectors;
import javax.sql.DataSource;

import org.apache.calcite.adapter.jdbc.JdbcSchema;
import org.apache.calcite.schema.Schema;
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,311 @@
package com.linkedin.hoptimator.util;

import java.util.Properties;
import java.sql.Array;
import java.sql.Blob;
import java.sql.CallableStatement;
import java.sql.Clob;
import java.sql.Connection;
import java.sql.DatabaseMetaData;
import java.sql.NClob;
import java.sql.PreparedStatement;
import java.sql.SQLClientInfoException;
import java.sql.SQLException;
import java.sql.SQLFeatureNotSupportedException;
import java.sql.SQLWarning;
import java.sql.SQLXML;
import java.sql.Savepoint;
import java.sql.Statement;
import java.sql.Struct;
import java.util.Map;
import java.util.concurrent.Executor;

class DelegatingConnection implements Connection {

private final Connection connection;

DelegatingConnection(Connection connection) {
this.connection = connection;
}

@Override
public String getSchema() throws SQLException {
return connection.getSchema();
}

@Override
public void setSchema(String schema) throws SQLException {
connection.setSchema(schema);
}

@Override
public String getCatalog() throws SQLException {
return connection.getCatalog();
}

@Override
public void setCatalog(String catalog) throws SQLException {
connection.setCatalog(catalog);
}

@Override
public void setAutoCommit(boolean autoCommit) throws SQLException {
// nop
}

@Override
public boolean getAutoCommit() throws SQLException {
return true;
}

@Override
public void close() throws SQLException {
connection.close();
}

@Override
public boolean isClosed() throws SQLException {
return connection.isClosed();
}

@Override
public DatabaseMetaData getMetaData() throws SQLException {
return connection.getMetaData();
}

@Override
public void setTransactionIsolation(int level) throws SQLException {
// nop
}

@Override
public int getTransactionIsolation() throws SQLException {
return Connection.TRANSACTION_NONE;
}

@Override
public SQLWarning getWarnings() throws SQLException {
return null;
}

@Override
public void setClientInfo(String name, String value) throws SQLClientInfoException {
connection.setClientInfo(name, value);
}

@Override
public void setClientInfo(Properties properties) throws SQLClientInfoException {
connection.setClientInfo(properties);
}

@Override
public String getClientInfo(String name) throws SQLException {
return connection.getClientInfo(name);
}

@Override
public Properties getClientInfo() throws SQLException {
return connection.getClientInfo();
}

@Override
public PreparedStatement prepareStatement(String sql) throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public CallableStatement prepareCall(String sql) throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public String nativeSQL(String sql) throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public void commit() throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public void rollback() throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public void setReadOnly(boolean readOnly) throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public boolean isReadOnly() throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public void clearWarnings() throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public Statement createStatement() throws SQLException {
return new DelegatingStatement(connection);
}

@Override
public Statement createStatement(int resultSetType, int resultSetConcurrency)
throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public PreparedStatement prepareStatement(
String sql, int resultSetType, int resultSetConcurrency) throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public CallableStatement prepareCall(String sql, int resultSetType, int resultSetConcurrency)
throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public Map<String, Class<?>> getTypeMap() throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public void setTypeMap(Map<String, Class<?>> map) throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public void setHoldability(int holdability) throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public int getHoldability() throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public Savepoint setSavepoint() throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public Savepoint setSavepoint(String name) throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public void rollback(Savepoint savepoint) throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public void releaseSavepoint(Savepoint savepoint) throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public Statement createStatement(
int resultSetType, int resultSetConcurrency, int resultSetHoldability)
throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public PreparedStatement prepareStatement(
String sql, int resultSetType, int resultSetConcurrency, int resultSetHoldability)
throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public CallableStatement prepareCall(
String sql, int resultSetType, int resultSetConcurrency, int resultSetHoldability)
throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public PreparedStatement prepareStatement(String sql, int autoGeneratedKeys)
throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public PreparedStatement prepareStatement(String sql, int[] columnIndexes) throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public PreparedStatement prepareStatement(String sql, String[] columnNames)
throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public Clob createClob() throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public Blob createBlob() throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public NClob createNClob() throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public SQLXML createSQLXML() throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public boolean isValid(int timeout) throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public Array createArrayOf(String typeName, Object[] elements) throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public Struct createStruct(String typeName, Object[] attributes) throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public void abort(Executor executor) throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public void setNetworkTimeout(Executor executor, int milliseconds) throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public int getNetworkTimeout() throws SQLException {
throw new SQLFeatureNotSupportedException();
}

@Override
public <T> T unwrap(Class<T> iface) throws SQLException {
return null;
}

@Override
public boolean isWrapperFor(Class<?> iface) throws SQLException {
return false;
}
}
Loading

0 comments on commit be410c9

Please sign in to comment.