Skip to content
Prev Previous commit
Next Next commit
Merge branch '1.10_release_4.0.x' into 1.10_release
# Conflicts:
#	core/src/main/java/com/dtstack/flink/sql/environment/MyLocalStreamEnvironment.java
#	core/src/main/java/com/dtstack/flink/sql/exec/ExecuteProcessHelper.java
#	hbase/hbase-side/hbase-async-side/src/main/java/com/dtstack/flink/sql/side/hbase/HbaseAsyncReqRow.java
  • Loading branch information
xuchao
xuchao committed Aug 2, 2020
commit cd71dfae984f980acc3a6c212ec36a961cb9076c
Original file line number Diff line number Diff line change
Expand Up @@ -106,13 +106,21 @@ public JobExecutionResult execute(StreamGraph streamGraph) throws Exception {
configuration.addAll(jobGraph.getJobConfiguration());

configuration.setString(TaskManagerOptions.MANAGED_MEMORY_SIZE.key(), "512M");
configuration.setInteger(TaskManagerOptions.NUM_TASK_SLOTS, jobGraph.getMaximumParallelism());

// add (and override) the settings with what the user defined
configuration.addAll(this.conf);

MiniClusterConfiguration.Builder configBuilder = new MiniClusterConfiguration.Builder();
configBuilder.setConfiguration(configuration);
configBuilder.setNumSlotsPerTaskManager(jobGraph.getMaximumParallelism());
if (!configuration.contains(RestOptions.BIND_PORT)) {
configuration.setString(RestOptions.BIND_PORT, "0");
}

int numSlotsPerTaskManager = configuration.getInteger(TaskManagerOptions.NUM_TASK_SLOTS, jobGraph.getMaximumParallelism());

MiniClusterConfiguration cfg = new MiniClusterConfiguration.Builder()
.setConfiguration(configuration)
.setNumSlotsPerTaskManager(numSlotsPerTaskManager)
.build();

if (LOG.isInfoEnabled()) {
LOG.info("Running job on local embedded Flink mini cluster");
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -62,15 +62,6 @@
import org.apache.calcite.sql.SqlNode;
import org.apache.commons.io.Charsets;
import org.apache.commons.lang3.StringUtils;
import org.apache.flink.api.common.typeinfo.TypeInformation;
import org.apache.flink.api.java.typeutils.RowTypeInfo;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.table.api.EnvironmentSettings;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.api.TableEnvironment;
import org.apache.flink.table.api.java.StreamTableEnvironment;
import org.apache.flink.table.sinks.TableSink;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

Expand All @@ -89,10 +80,9 @@
import java.util.ArrayList;

/**
* 任务执行时的流程方法
* 任务执行时的流程方法
* Date: 2020/2/17
* Company: www.dtstack.com
*
* @author maqi
*/
public class ExecuteProcessHelper {
Expand Down Expand Up @@ -137,11 +127,11 @@ public static ParamsInfo parseParams(String[] args) throws Exception {
.setConfProp(confProperties)
.setJarUrlList(jarUrlList)
.build();

}

/**
* 非local模式或者shipfile部署模式,remoteSqlPluginPath必填
*
* 非local模式或者shipfile部署模式,remoteSqlPluginPath必填
* @param remoteSqlPluginPath
* @param deployMode
* @param pluginLoadMode
Expand All @@ -160,6 +150,7 @@ public static StreamExecutionEnvironment getStreamExecution(ParamsInfo paramsInf
StreamExecutionEnvironment env = ExecuteProcessHelper.getStreamExeEnv(paramsInfo.getConfProp(), paramsInfo.getDeployMode());
StreamTableEnvironment tableEnv = getStreamTableEnv(env, paramsInfo.getConfProp());


SqlParser.setLocalSqlPluginRoot(paramsInfo.getLocalSqlPluginPath());
SqlTree sqlTree = SqlParser.parseSql(paramsInfo.getSql(), paramsInfo.getPluginLoadMode());

Expand Down Expand Up @@ -200,7 +191,7 @@ public static List<URL> getExternalJarUrls(String addJarListStr) throws java.io.
private static void sqlTranslation(String localSqlPluginPath,
String pluginLoadMode,
StreamTableEnvironment tableEnv,
SqlTree sqlTree, Map<String, AbstractSideTableInfo> sideTableMap,
SqlTree sqlTree,Map<String, AbstractSideTableInfo> sideTableMap,
Map<String, Table> registerTableCache) throws Exception {

SideSqlExec sideSqlExec = new SideSqlExec();
Expand Down Expand Up @@ -268,14 +259,13 @@ public static void registerUserDefinedFunction(SqlTree sqlTree, List<URL> jarUrl
}

/**
* 向Flink注册源表和结果表,返回执行时插件包的全路径
*
* 向Flink注册源表和结果表,返回执行时插件包的全路径
* @param sqlTree
* @param env
* @param tableEnv
* @param localSqlPluginPath
* @param remoteSqlPluginPath
* @param pluginLoadMode 插件加载模式 classpath or shipfile
* @param pluginLoadMode 插件加载模式 classpath or shipfile
* @param sideTableMap
* @param registerTableCache
* @return
Expand Down Expand Up @@ -343,8 +333,7 @@ public static Set<URL> registerTable(SqlTree sqlTree, StreamExecutionEnvironment
}

/**
* perjob模式将job依赖的插件包路径存储到cacheFile,在外围将插件包路径传递给jobgraph
*
* perjob模式将job依赖的插件包路径存储到cacheFile,在外围将插件包路径传递给jobgraph
* @param env
* @param classPathSet
*/
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,8 @@
* limitations under the License.
*/



package com.dtstack.flink.sql.side.hbase;

import com.dtstack.flink.sql.enums.ECacheContentType;
Expand Down Expand Up @@ -45,7 +47,12 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.*;
import java.io.File;
import java.io.IOException;
import java.sql.Timestamp;
import java.util.Collections;
import java.util.List;
import java.util.Map;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
Expand Down Expand Up @@ -74,9 +81,9 @@ public class HbaseAsyncReqRow extends BaseAsyncReqRow {

private transient AbstractRowKeyModeDealer rowKeyMode;

private final String tableName;
private String tableName;

private final String[] colNames;
private String[] colNames;

public HbaseAsyncReqRow(RowTypeInfo rowTypeInfo, JoinInfo joinInfo, List<FieldInfo> outFieldInfoList, AbstractSideTableInfo sideTableInfo) {
super(new HbaseAsyncSideInfo(rowTypeInfo, joinInfo, outFieldInfoList, sideTableInfo));
Expand Down Expand Up @@ -139,18 +146,9 @@ public void open(Configuration parameters) throws Exception {
}

@Override
public void asyncInvoke(Tuple2<Boolean,Row> input, ResultFuture<Tuple2<Boolean,Row>> resultFuture) throws Exception {
Tuple2<Boolean,Row> inputCopy = Tuple2.of(input.f0,input.f1);
Map<String, Object> refData = Maps.newHashMap();
for (int i = 0; i < sideInfo.getEqualValIndex().size(); i++) {
Integer conValIndex = sideInfo.getEqualValIndex().get(i);
Object equalObj = inputCopy.f1.getField(conValIndex);
if(equalObj == null){
dealMissKey(inputCopy, resultFuture);
return;
}
refData.put(getAliasFieldsName(sideInfo.getEqualFieldList().get(i), sideInfo.getSideTableInfo().getPhysicalFields()), equalObj);
}
public void handleAsyncInvoke(Map<String, Object> inputParams, Row input, ResultFuture<BaseRow> resultFuture) throws Exception {
rowKeyMode.asyncGetData(tableName, buildCacheKey(inputParams), input, resultFuture, sideInfo.getSideCache());
}

@Override
public String buildCacheKey(Map<String, Object> inputParams) {
Expand Down Expand Up @@ -178,23 +176,6 @@ public Row fillData(Row input, Object sideInput){
return row;
}

// 根据实际字段名获得对应的别名
public String getAliasFieldsName(String realFieldName, Map<String, String> physicalFields) {
Collection<String> values = physicalFields.values();
Set<String> keySet = physicalFields.keySet();
if (!values.contains(realFieldName)) {
// TODO Error ? or Warn ?
LOG.warn(realFieldName + "不存在别名");
} else {
for (String key : keySet) {
if (physicalFields.get(key).equals(realFieldName)) {
return key;
}
}
}
return realFieldName;
}

@Override
public void close() throws Exception {
super.close();
Expand Down
You are viewing a condensed version of this merge commit. You can view the full changes here.