Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
73 changes: 73 additions & 0 deletions connectors/connector-pravega/pom.xml
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<parent>
<artifactId>alink_connectors</artifactId>
<groupId>com.alibaba.alink</groupId>
<version>1.5-SNAPSHOT</version>
</parent>
<modelVersion>4.0.0</modelVersion>

<artifactId>alink_connector_pravega_flink-${alink.flink.major.version}_${alink.scala.major.version}</artifactId>
<name>alink-connector-pravega</name>

<packaging>jar</packaging>

<build>
<plugins>
<plugin>
<groupId>org.apache.maven.plugins</groupId>
<artifactId>maven-compiler-plugin</artifactId>
<configuration>
<source>1.8</source>
<target>1.8</target>
</configuration>
</plugin>
</plugins>
</build>

<dependencies>
<dependency>
<groupId>com.alibaba.alink</groupId>
<artifactId>alink_core_flink-${alink.flink.major.version}_${alink.scala.major.version}</artifactId>
<version>${project.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-java_${alink.scala.major.version}</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-streaming-scala_${alink.scala.major.version}</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-api-java-bridge_${alink.scala.major.version}</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-api-java</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>org.apache.flink</groupId>
<artifactId>flink-table-planner_${alink.scala.major.version}</artifactId>
<version>${flink.version}</version>
<scope>provided</scope>
</dependency>
<dependency>
<groupId>com.jd.flink.sql</groupId>
<artifactId>flink-connector-pravega</artifactId>
<version>1.0-SNAPSHOT</version>
</dependency>
</dependencies>
</project>
Original file line number Diff line number Diff line change
@@ -0,0 +1,96 @@
package com.alibaba.alink.common.io.pravega.plugin;

import com.alibaba.alink.operator.common.io.serde.RowToCsvSerialization;
import io.pravega.client.admin.StreamManager;
import io.pravega.client.stream.Stream;
import io.pravega.client.stream.StreamConfiguration;
import io.pravega.connectors.flink.FlinkPravegaInputFormat;
import io.pravega.connectors.flink.FlinkPravegaReader;
import io.pravega.connectors.flink.FlinkPravegaWriter;
import io.pravega.connectors.flink.PravegaConfig;
import org.apache.flink.api.common.restartstrategy.RestartStrategies;
import org.apache.flink.api.common.serialization.SimpleStringSchema;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.java.tuple.Tuple1;
import org.apache.flink.api.java.utils.ParameterTool;
import org.apache.flink.ml.api.misc.param.Params;
import org.apache.flink.streaming.api.CheckpointingMode;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;
import org.apache.flink.streaming.api.functions.source.RichParallelSourceFunction;
import org.apache.flink.types.Row;
import java.net.URI;

public class PravegaSourceSinkInPluginFactory implements PravegaSourceSinkFactory {


@Override
public Tuple1<RichParallelSourceFunction> createPravegaSourceFunction(Params params) {
String DEFAULT_SCOPE = params.getStringOrDefault("scope", "scope");
String Default_URI_PARAM = params.getStringOrDefault("controller", "controller");
String STREAM_PARAM = params.getStringOrDefault("stream", "stream");

PravegaConfig pravegaConfig = PravegaConfig
.fromParams(null)
.withControllerURI(URI.create(Default_URI_PARAM))
.withDefaultScope(DEFAULT_SCOPE);
Stream stream = pravegaConfig.resolve(STREAM_PARAM);
try(StreamManager streamManager = StreamManager.create(pravegaConfig.getClientConfig())) {
streamManager.createScope(stream.getScope());
streamManager.createStream(stream.getScope(), stream.getStreamName(), StreamConfiguration.builder().build());
}
FlinkPravegaReader<String> source = FlinkPravegaReader.<String>builder()
.withPravegaConfig(pravegaConfig)
.forStream(stream)
.withDeserializationSchema(new SimpleStringSchema())
.build();
return Tuple1.of(source);
}


@Override
public Tuple1<FlinkPravegaInputFormat> createPravegaSourceFunctionBatch(Params params) {
String DEFAULT_SCOPE = params.getStringOrDefault("scope", "scope");
String Default_URI_PARAM = params.getStringOrDefault("controller", "controller");
String STREAM_PARAM = params.getStringOrDefault("stream", "stream");
PravegaConfig pravegaConfig = PravegaConfig
.fromParams(null)
.withControllerURI(URI.create(Default_URI_PARAM))
.withDefaultScope(DEFAULT_SCOPE);
Stream stream = pravegaConfig.resolve(STREAM_PARAM);
try(StreamManager streamManager = StreamManager.create(pravegaConfig.getClientConfig())) {
streamManager.createScope(stream.getScope());
streamManager.createStream(stream.getScope(), stream.getStreamName(), StreamConfiguration.builder().build());
}
FlinkPravegaInputFormat<String> source = FlinkPravegaInputFormat.<String>builder()
.forStream(stream)
.withPravegaConfig(pravegaConfig)
.withDeserializationSchema(new SimpleStringSchema())
.build();
return Tuple1.of(source);
}

@Override
public Tuple1<RichSinkFunction> createPravegaSinkFunction(Params params) {
String DEFAULT_SCOPE = params.getStringOrDefault("scope", "scope");
String Default_URI_PARAM = params.getStringOrDefault("controller", "controller");
String STREAM_PARAM = params.getStringOrDefault("stream", "stream");
String fieldDelim = params.getStringOrDefault("fieldDelim", ",");
PravegaConfig pravegaConfig = PravegaConfig
.fromParams(null)
.withControllerURI(URI.create(Default_URI_PARAM))
.withDefaultScope(DEFAULT_SCOPE);
Stream stream = pravegaConfig.resolve(STREAM_PARAM);
try(StreamManager streamManager = StreamManager.create(pravegaConfig.getClientConfig())) {
streamManager.createScope(stream.getScope());
streamManager.createStream(stream.getScope(), stream.getStreamName(), StreamConfiguration.builder().build());
FlinkPravegaWriter<Row> writer = FlinkPravegaWriter.<Row>builder()
.withPravegaConfig(pravegaConfig)
.forStream(stream)
.withSerializationSchema(new RowToCsvSerialization(new TypeInformation[]{TypeInformation.of(String.class)}, fieldDelim))
.build();
return Tuple1.of(writer);
}

}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
com.alibaba.alink.common.io.pravega.plugin.PravegaSourceSinkInPluginFactory
1 change: 1 addition & 0 deletions connectors/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
<module>connector-odps</module>
<module>connector-jdbc</module>
<module>connector-hive</module>
<module>connector-pravega</module>
<module>filesystem</module>
</modules>

Expand Down
11 changes: 11 additions & 0 deletions core/pom.xml
Original file line number Diff line number Diff line change
Expand Up @@ -248,6 +248,12 @@
<version>0.9.10</version>
</dependency>

<dependency>
<groupId>net.sf.py4j</groupId>
<artifactId>py4j</artifactId>
<version>0.10.9.2</version>
</dependency>

<dependency>
<groupId>com.github.fommil.netlib</groupId>
<artifactId>all</artifactId>
Expand Down Expand Up @@ -323,6 +329,11 @@
<type>jar</type>
<scope>test</scope>
</dependency>
<dependency>
<groupId>com.jd.flink.sql</groupId>
<artifactId>flink-connector-pravega</artifactId>
<version>1.0-SNAPSHOT</version>
</dependency>
</dependencies>

</project>
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
package com.alibaba.alink.common.io.pravega.plugin;

import com.alibaba.alink.common.io.plugin.*;
import org.apache.flink.api.java.tuple.Tuple2;

import java.lang.reflect.InvocationTargetException;
import java.util.function.Function;
import java.util.function.Predicate;

public class PravegaClassLoaderFactory extends ClassLoaderFactory {
public final static String PRAVEGA_NAME = "pravega";

public PravegaClassLoaderFactory(String version) {
super(new RegisterKey(PRAVEGA_NAME, version), PluginDistributeCache.createDistributeCache(PRAVEGA_NAME, version));
}

public static PravegaSourceSinkFactory create(PravegaClassLoaderFactory factory) {
try {
return (PravegaSourceSinkFactory) factory
.create()
.loadClass("com.alibaba.alink.common.io.pravega.plugin.PravegaSourceSinkInPluginFactory")
.getConstructor()
.newInstance();
} catch (InstantiationException | IllegalAccessException | InvocationTargetException | NoSuchMethodException |
ClassNotFoundException e) {
throw new RuntimeException(e);
}
}

@Override
public ClassLoader create() {
return ClassLoaderContainer.getInstance().create(
registerKey,
distributeCache,
PravegaSourceSinkFactory.class,
new PravegaServiceFilter(),
new PravegaVersionGetter()
);
}

private static class PravegaServiceFilter implements Predicate <PravegaSourceSinkFactory> {

@Override
public boolean test(PravegaSourceSinkFactory factory) {
return true;
}
}

private static class PravegaVersionGetter implements
Function <Tuple2<PravegaSourceSinkFactory, PluginDescriptor>, String> {

@Override
public String apply(Tuple2<PravegaSourceSinkFactory, PluginDescriptor> factory) {
return "0001";
}
}

}
Original file line number Diff line number Diff line change
@@ -0,0 +1,13 @@
package com.alibaba.alink.common.io.pravega.plugin;

import io.pravega.connectors.flink.FlinkPravegaInputFormat;
import org.apache.flink.api.java.tuple.Tuple1;
import org.apache.flink.ml.api.misc.param.Params;
import org.apache.flink.streaming.api.functions.sink.RichSinkFunction;
import org.apache.flink.streaming.api.functions.source.RichParallelSourceFunction;

public interface PravegaSourceSinkFactory {
Tuple1<RichParallelSourceFunction> createPravegaSourceFunction(Params params);
Tuple1<FlinkPravegaInputFormat> createPravegaSourceFunctionBatch(Params params);
Tuple1<RichSinkFunction> createPravegaSinkFunction(Params params);
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
package com.alibaba.alink.operator.batch.source;

import com.alibaba.alink.common.MLEnvironmentFactory;
import com.alibaba.alink.common.io.annotations.AnnotationUtils;
import com.alibaba.alink.common.io.annotations.IOType;
import com.alibaba.alink.common.io.annotations.IoOpAnnotation;
import com.alibaba.alink.common.io.pravega.plugin.PravegaClassLoaderFactory;
import com.alibaba.alink.common.utils.DataStreamConversionUtil;
import com.alibaba.alink.operator.stream.source.BaseSourceStreamOp;
import com.alibaba.alink.params.io.KafkaSourceParams;
import org.apache.flink.api.java.tuple.Tuple1;
import org.apache.flink.ml.api.misc.param.Params;
import org.apache.flink.streaming.api.functions.source.RichParallelSourceFunction;
import org.apache.flink.table.api.Table;

@IoOpAnnotation(name = "pravega", ioType = IOType.SourceStream)
public class PravegaSourceBatchOp extends BaseSourceStreamOp<PravegaSourceBatchOp>
implements KafkaSourceParams<PravegaSourceBatchOp> {

private PravegaClassLoaderFactory factory;

public PravegaSourceBatchOp() {
this(new Params());
}

public PravegaSourceBatchOp(Params params) {
super(AnnotationUtils.annotatedName(PravegaSourceBatchOp.class), params);
}

@Override
protected Table initializeDataSource() {
if (factory == null) {
factory = new PravegaClassLoaderFactory("0001");
}

Tuple1<RichParallelSourceFunction> sourceFunction = PravegaClassLoaderFactory
.create(factory)
.createPravegaSourceFunction(getParams());

return DataStreamConversionUtil.toTable(
getMLEnvironmentId(),
MLEnvironmentFactory.get(getMLEnvironmentId())
.getStreamExecutionEnvironment()
.addSource(sourceFunction.f0),
new String[]{}
);
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,51 @@
package com.alibaba.alink.operator.common.prophet;

import org.apache.flink.api.common.functions.RichMapPartitionFunction;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.types.Row;
import org.apache.flink.util.Collector;

import java.util.Iterator;

/**
* the way to build the prophet model
*/
public class BuildProphetModel extends RichMapPartitionFunction<Row, Row> {
private static final long serialVersionUID = -4019434236075549258L;
private String[] selectedColNames;

public BuildProphetModel(String[] selectedColNames) {
this.selectedColNames = selectedColNames;
}


@Override
public void open(Configuration parameters) throws Exception {
super.open(parameters);
}

@Override
public void mapPartition(Iterable<Row> iterable, Collector<Row> collector) throws Exception {
StringBuilder sbd = new StringBuilder();
Iterator<Row> iter = iterable.iterator();
while(iter.hasNext()) {
Row row = iter.next();
sbd.append(row.getField(0)).append(",");
sbd.append(row.getField(1)).append(";");
}
ListenerApplication.open();
ListenerApplication.modelParam = sbd.toString().substring(0, sbd.length() - 1);
Runtime rt = Runtime.getRuntime();
Process process = rt.exec("python " + BuildProphetModel.class.getClassLoader().getResource("python_file/prophet/prophet_demo.py").getPath());
process.waitFor();
ListenerApplication.close();
Row row = new Row(1);
row.setField(0, ListenerApplication.returnValue);
collector.collect(row);
}

@Override
public void close() throws Exception {
super.close();
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,8 @@
package com.alibaba.alink.operator.common.prophet;

/**
* use py4j,the interface that python code implements
*/
public interface ExampleListener {
public Object notify(Object source);
}
Loading