Skip to content
Projects
Groups
Snippets
Help
Loading...
Help
Contribute to GitLab
Sign in / Register
Toggle navigation
D
dlink
Project
Project
Details
Activity
Cycle Analytics
Repository
Repository
Files
Commits
Branches
Tags
Contributors
Graph
Compare
Charts
Issues
0
Issues
0
List
Board
Labels
Milestones
Merge Requests
0
Merge Requests
0
CI / CD
CI / CD
Pipelines
Jobs
Schedules
Charts
Wiki
Wiki
Snippets
Snippets
Members
Members
Collapse sidebar
Close sidebar
Activity
Graph
Charts
Create a new issue
Jobs
Commits
Issue Boards
Open sidebar
zhaowei
dlink
Commits
be150742
Commit
be150742
authored
Mar 09, 2022
by
wenmo
Browse files
Options
Browse Files
Download
Email Patches
Plain Diff
优化K8S Application提交配置
parent
a68a902e
Changes
12
Show whitespace changes
Inline
Side-by-side
Showing
12 changed files
with
72 additions
and
151 deletions
+72
-151
MainApp.java
dlink-app/src/main/java/com/dlink/app/MainApp.java
+9
-8
DBConfig.java
dlink-app/src/main/java/com/dlink/app/db/DBConfig.java
+8
-6
CustomTableEnvironmentImpl.java
...n/java/com/dlink/executor/CustomTableEnvironmentImpl.java
+1
-1
CustomTableResultImpl.java
...c/main/java/com/dlink/executor/CustomTableResultImpl.java
+1
-1
pom.xml
dlink-client/dlink-client-base/pom.xml
+0
-124
FlinkParamConstant.java
.../src/main/java/com/dlink/constant/FlinkParamConstant.java
+15
-0
FlinkBaseUtil.java
...ent-base/src/main/java/com/dlink/utils/FlinkBaseUtil.java
+27
-0
pom.xml
dlink-executor/pom.xml
+0
-8
KubernetesGateway.java
.../java/com/dlink/gateway/kubernetes/KubernetesGateway.java
+1
-2
ClusterConfigurationForm.tsx
...sterConfiguration/components/ClusterConfigurationForm.tsx
+1
-1
conf.ts
dlink-web/src/pages/ClusterConfiguration/conf.ts
+6
-0
Welcome.tsx
dlink-web/src/pages/Welcome.tsx
+3
-0
No files found.
dlink-app/src/main/java/com/dlink/app/MainApp.java
View file @
be150742
...
@@ -2,10 +2,12 @@ package com.dlink.app;
...
@@ -2,10 +2,12 @@ package com.dlink.app;
import
com.dlink.app.db.DBConfig
;
import
com.dlink.app.db.DBConfig
;
import
com.dlink.app.flinksql.Submiter
;
import
com.dlink.app.flinksql.Submiter
;
import
org.apache.flink.api.java.utils.ParameterTool
;
import
com.dlink.assertion.Asserts
;
import
com.dlink.constant.FlinkParamConstant
;
import
com.dlink.utils.FlinkBaseUtil
;
import
java.io.IOException
;
import
java.io.IOException
;
import
java.
time.LocalDateTime
;
import
java.
util.Map
;
/**
/**
* MainApp
* MainApp
...
@@ -16,11 +18,10 @@ import java.time.LocalDateTime;
...
@@ -16,11 +18,10 @@ import java.time.LocalDateTime;
public
class
MainApp
{
public
class
MainApp
{
public
static
void
main
(
String
[]
args
)
throws
IOException
{
public
static
void
main
(
String
[]
args
)
throws
IOException
{
ParameterTool
parameters
=
ParameterTool
.
fromArgs
(
args
);
Map
<
String
,
String
>
params
=
FlinkBaseUtil
.
getParamsFromArgs
(
args
);
String
id
=
parameters
.
get
(
"id"
,
null
);
String
id
=
params
.
get
(
FlinkParamConstant
.
ID
);
if
(
id
!=
null
&&!
""
.
equals
(
id
))
{
Asserts
.
checkNullString
(
id
,
"请配置入参 id "
);
DBConfig
dbConfig
=
DBConfig
.
build
(
parameters
);
DBConfig
dbConfig
=
DBConfig
.
build
(
params
);
Submiter
.
submit
(
Integer
.
valueOf
(
id
),
dbConfig
);
Submiter
.
submit
(
Integer
.
valueOf
(
id
),
dbConfig
);
}
}
}
}
}
dlink-app/src/main/java/com/dlink/app/db/DBConfig.java
View file @
be150742
package
com
.
dlink
.
app
.
db
;
package
com
.
dlink
.
app
.
db
;
import
org.apache.flink.api.java.utils.ParameterTool
;
import
com.dlink.constant.FlinkParamConstant
;
import
java.util.Map
;
/**
/**
* DBConfig
* DBConfig
...
@@ -27,11 +29,11 @@ public class DBConfig {
...
@@ -27,11 +29,11 @@ public class DBConfig {
}
}
public
static
DBConfig
build
(
ParameterTool
parameter
s
){
public
static
DBConfig
build
(
Map
<
String
,
String
>
param
s
){
return
new
DBConfig
(
param
eters
.
get
(
"driver"
,
null
),
return
new
DBConfig
(
param
s
.
get
(
FlinkParamConstant
.
DRIVER
),
param
eters
.
get
(
"url"
,
null
),
param
s
.
get
(
FlinkParamConstant
.
URL
),
param
eters
.
get
(
"username"
,
null
),
param
s
.
get
(
FlinkParamConstant
.
USERNAME
),
param
eters
.
get
(
"password"
,
null
));
param
s
.
get
(
FlinkParamConstant
.
PASSWORD
));
}
}
public
String
getDriver
()
{
public
String
getDriver
()
{
...
...
dlink-app/src/main/java/com/dlink/executor/CustomTableEnvironmentImpl.java
View file @
be150742
...
@@ -55,7 +55,7 @@ import java.util.Map;
...
@@ -55,7 +55,7 @@ import java.util.Map;
* @author wenmo
* @author wenmo
* @since 2021/6/7 22:06
* @since 2021/6/7 22:06
**/
**/
public
class
CustomTableEnvironmentImpl
extends
TableEnvironmentImpl
implements
CustomTableEnvironment
{
public
class
CustomTableEnvironmentImpl
extends
TableEnvironmentImpl
implements
CustomTableEnvironment
{
protected
CustomTableEnvironmentImpl
(
CatalogManager
catalogManager
,
ModuleManager
moduleManager
,
TableConfig
tableConfig
,
Executor
executor
,
FunctionCatalog
functionCatalog
,
Planner
planner
,
boolean
isStreamingMode
,
ClassLoader
userClassLoader
)
{
protected
CustomTableEnvironmentImpl
(
CatalogManager
catalogManager
,
ModuleManager
moduleManager
,
TableConfig
tableConfig
,
Executor
executor
,
FunctionCatalog
functionCatalog
,
Planner
planner
,
boolean
isStreamingMode
,
ClassLoader
userClassLoader
)
{
super
(
catalogManager
,
moduleManager
,
tableConfig
,
executor
,
functionCatalog
,
planner
,
isStreamingMode
,
userClassLoader
);
super
(
catalogManager
,
moduleManager
,
tableConfig
,
executor
,
functionCatalog
,
planner
,
isStreamingMode
,
userClassLoader
);
...
...
dlink-app/src/main/java/com/dlink/executor/CustomTableResultImpl.java
View file @
be150742
...
@@ -59,7 +59,7 @@ public class CustomTableResultImpl implements TableResult {
...
@@ -59,7 +59,7 @@ public class CustomTableResultImpl implements TableResult {
Preconditions
.
checkNotNull
(
sessionTimeZone
,
"sessionTimeZone should not be null"
);
Preconditions
.
checkNotNull
(
sessionTimeZone
,
"sessionTimeZone should not be null"
);
}
}
public
static
TableResult
buildTableResult
(
List
<
TableSchemaField
>
fields
,
List
<
Row
>
rows
){
public
static
TableResult
buildTableResult
(
List
<
TableSchemaField
>
fields
,
List
<
Row
>
rows
){
Builder
builder
=
builder
().
resultKind
(
ResultKind
.
SUCCESS
);
Builder
builder
=
builder
().
resultKind
(
ResultKind
.
SUCCESS
);
if
(
fields
.
size
()>
0
)
{
if
(
fields
.
size
()>
0
)
{
List
<
String
>
columnNames
=
new
ArrayList
<>();
List
<
String
>
columnNames
=
new
ArrayList
<>();
...
...
dlink-client/dlink-client-base/pom.xml
View file @
be150742
...
@@ -24,39 +24,6 @@
...
@@ -24,39 +24,6 @@
<maven.compiler.target>
8
</maven.compiler.target>
<maven.compiler.target>
8
</maven.compiler.target>
</properties>
</properties>
<dependencyManagement>
<dependencies>
<dependency>
<groupId>
org.springframework
</groupId>
<artifactId>
spring-framework-bom
</artifactId>
<version>
${spring.version}
</version>
<type>
pom
</type>
<scope>
import
</scope>
</dependency>
<dependency>
<groupId>
org.apache.dubbo
</groupId>
<artifactId>
dubbo-bom
</artifactId>
<version>
${dubbo.version}
</version>
<type>
pom
</type>
<scope>
import
</scope>
</dependency>
<dependency>
<groupId>
org.apache.dubbo
</groupId>
<artifactId>
dubbo-dependencies-zookeeper
</artifactId>
<version>
${dubbo.version}
</version>
<type>
pom
</type>
</dependency>
<dependency>
<groupId>
junit
</groupId>
<artifactId>
junit
</artifactId>
<version>
${junit.version}
</version>
</dependency>
</dependencies>
</dependencyManagement>
<dependencies>
<dependencies>
<dependency>
<dependency>
<groupId>
com.dlink
</groupId>
<groupId>
com.dlink
</groupId>
...
@@ -98,97 +65,7 @@
...
@@ -98,97 +65,7 @@
<version>
${flink.version}
</version>
<version>
${flink.version}
</version>
<scope>
provided
</scope>
<scope>
provided
</scope>
</dependency>
</dependency>
<dependency>
<groupId>
org.apache.dubbo
</groupId>
<artifactId>
dubbo
</artifactId>
<scope>
provided
</scope>
</dependency>
<dependency>
<groupId>
org.apache.dubbo
</groupId>
<artifactId>
dubbo-dependencies-zookeeper
</artifactId>
<type>
pom
</type>
<scope>
provided
</scope>
</dependency>
<dependency>
<groupId>
junit
</groupId>
<artifactId>
junit
</artifactId>
<scope>
test
</scope>
</dependency>
<dependency>
<groupId>
org.springframework
</groupId>
<artifactId>
spring-test
</artifactId>
<scope>
test
</scope>
</dependency>
</dependencies>
</dependencies>
<profiles>
<!-- For jdk 11 above JavaEE annotation -->
<profile>
<id>
javax.annotation
</id>
<activation>
<jdk>
[1.11,)
</jdk>
</activation>
<dependencies>
<dependency>
<groupId>
javax.annotation
</groupId>
<artifactId>
javax.annotation-api
</artifactId>
<version>
1.3.2
</version>
</dependency>
</dependencies>
</profile>
<profile>
<id>
provider
</id>
<build>
<plugins>
<plugin>
<!-- Build an executable JAR -->
<groupId>
org.apache.maven.plugins
</groupId>
<artifactId>
maven-jar-plugin
</artifactId>
<version>
3.1.0
</version>
<configuration>
<finalName>
provider
</finalName>
<archive>
<manifest>
<addClasspath>
true
</addClasspath>
<classpathPrefix>
lib/
</classpathPrefix>
<mainClass>
com.dlink.BasicProvider
</mainClass>
</manifest>
</archive>
</configuration>
</plugin>
</plugins>
</build>
</profile>
<profile>
<id>
consumer
</id>
<build>
<plugins>
<plugin>
<!-- Build an executable JAR -->
<groupId>
org.apache.maven.plugins
</groupId>
<artifactId>
maven-jar-plugin
</artifactId>
<version>
3.1.0
</version>
<configuration>
<finalName>
consumer
</finalName>
<archive>
<manifest>
<addClasspath>
true
</addClasspath>
<classpathPrefix>
lib/
</classpathPrefix>
<mainClass>
com.dlink.BasicConsumer
</mainClass>
</manifest>
</archive>
</configuration>
</plugin>
</plugins>
</build>
</profile>
</profiles>
<build>
<build>
<plugins>
<plugins>
<plugin>
<plugin>
...
@@ -200,7 +77,6 @@
...
@@ -200,7 +77,6 @@
<target>
${target.level}
</target>
<target>
${target.level}
</target>
</configuration>
</configuration>
</plugin>
</plugin>
<plugin>
<plugin>
<groupId>
org.apache.maven.plugins
</groupId>
<groupId>
org.apache.maven.plugins
</groupId>
<artifactId>
maven-dependency-plugin
</artifactId>
<artifactId>
maven-dependency-plugin
</artifactId>
...
...
dlink-client/dlink-client-base/src/main/java/com/dlink/constant/FlinkParamConstant.java
0 → 100644
View file @
be150742
package
com
.
dlink
.
constant
;
/**
* FlinkParam
*
* @author wenmo
* @since 2022/3/9 19:18
*/
public
final
class
FlinkParamConstant
{
public
static
final
String
ID
=
"id"
;
public
static
final
String
DRIVER
=
"driver"
;
public
static
final
String
URL
=
"url"
;
public
static
final
String
USERNAME
=
"username"
;
public
static
final
String
PASSWORD
=
"password"
;
}
dlink-client/dlink-client-base/src/main/java/com/dlink/utils/FlinkBaseUtil.java
0 → 100644
View file @
be150742
package
com
.
dlink
.
utils
;
import
com.dlink.constant.FlinkParamConstant
;
import
org.apache.flink.api.java.utils.ParameterTool
;
import
java.util.HashMap
;
import
java.util.Map
;
/**
* FlinkBaseUtil
*
* @author wenmo
* @since 2022/3/9 19:15
*/
public
class
FlinkBaseUtil
{
public
static
Map
<
String
,
String
>
getParamsFromArgs
(
String
[]
args
){
Map
<
String
,
String
>
params
=
new
HashMap
<>();
ParameterTool
parameters
=
ParameterTool
.
fromArgs
(
args
);
params
.
put
(
FlinkParamConstant
.
ID
,
parameters
.
get
(
FlinkParamConstant
.
ID
,
null
));
params
.
put
(
FlinkParamConstant
.
DRIVER
,
parameters
.
get
(
FlinkParamConstant
.
DRIVER
,
null
));
params
.
put
(
FlinkParamConstant
.
URL
,
parameters
.
get
(
FlinkParamConstant
.
URL
,
null
));
params
.
put
(
FlinkParamConstant
.
USERNAME
,
parameters
.
get
(
FlinkParamConstant
.
USERNAME
,
null
));
params
.
put
(
FlinkParamConstant
.
PASSWORD
,
parameters
.
get
(
FlinkParamConstant
.
PASSWORD
,
null
));
return
params
;
}
}
dlink-executor/pom.xml
View file @
be150742
...
@@ -22,14 +22,6 @@
...
@@ -22,14 +22,6 @@
<groupId>
com.dlink
</groupId>
<groupId>
com.dlink
</groupId>
<artifactId>
dlink-common
</artifactId>
<artifactId>
dlink-common
</artifactId>
</dependency>
</dependency>
<dependency>
<groupId>
com.fasterxml.jackson.core
</groupId>
<artifactId>
jackson-annotations
</artifactId>
</dependency>
<dependency>
<groupId>
com.fasterxml.jackson.core
</groupId>
<artifactId>
jackson-databind
</artifactId>
</dependency>
<dependency>
<dependency>
<groupId>
junit
</groupId>
<groupId>
junit
</groupId>
<artifactId>
junit
</artifactId>
<artifactId>
junit
</artifactId>
...
...
dlink-gateway/src/main/java/com/dlink/gateway/kubernetes/KubernetesGateway.java
View file @
be150742
...
@@ -57,9 +57,8 @@ public abstract class KubernetesGateway extends AbstractGateway {
...
@@ -57,9 +57,8 @@ public abstract class KubernetesGateway extends AbstractGateway {
if
(
Asserts
.
isNotNullString
(
config
.
getFlinkConfig
().
getSavePoint
()))
{
if
(
Asserts
.
isNotNullString
(
config
.
getFlinkConfig
().
getSavePoint
()))
{
configuration
.
setString
(
SavepointConfigOptions
.
SAVEPOINT_PATH
,
config
.
getFlinkConfig
().
getSavePoint
());
configuration
.
setString
(
SavepointConfigOptions
.
SAVEPOINT_PATH
,
config
.
getFlinkConfig
().
getSavePoint
());
}
}
configuration
.
set
(
YarnConfigOptions
.
PROVIDED_LIB_DIRS
,
Collections
.
singletonList
(
config
.
getClusterConfig
().
getFlinkLibPath
()));
if
(
Asserts
.
isNotNullString
(
config
.
getFlinkConfig
().
getJobName
()))
{
if
(
Asserts
.
isNotNullString
(
config
.
getFlinkConfig
().
getJobName
()))
{
configuration
.
set
(
YarnConfigOptions
.
APPLICATION_NAME
,
config
.
getFlinkConfig
().
getJobName
());
configuration
.
set
(
KubernetesConfigOptions
.
CLUSTER_ID
,
config
.
getFlinkConfig
().
getJobName
());
}
}
}
}
...
...
dlink-web/src/pages/ClusterConfiguration/components/ClusterConfigurationForm.tsx
View file @
be150742
...
@@ -50,7 +50,7 @@ const ClusterConfigurationForm: React.FC<ClusterConfigurationFormProps> = (props
...
@@ -50,7 +50,7 @@ const ClusterConfigurationForm: React.FC<ClusterConfigurationFormProps> = (props
name=
{
configItem
.
name
}
name=
{
configItem
.
name
}
label=
{
configItem
.
lable
}
label=
{
configItem
.
lable
}
>
>
<
Input
placeholder=
{
configItem
.
placeholder
}
/>
<
Input
placeholder=
{
configItem
.
placeholder
}
defaultValue=
{
configItem
.
defaultValue
}
/>
</
Form
.
Item
>)
</
Form
.
Item
>)
});
});
return
itemList
;
return
itemList
;
...
...
dlink-web/src/pages/ClusterConfiguration/conf.ts
View file @
be150742
...
@@ -2,6 +2,7 @@ export type Config = {
...
@@ -2,6 +2,7 @@ export type Config = {
name
:
string
,
name
:
string
,
lable
:
string
,
lable
:
string
,
placeholder
:
string
placeholder
:
string
defaultValue
?:
string
}
}
export
const
HADOOP_CONFIG_LIST
:
Config
[]
=
[{
export
const
HADOOP_CONFIG_LIST
:
Config
[]
=
[{
...
@@ -21,6 +22,11 @@ export const KUBERNETES_CONFIG_LIST: Config[] = [{
...
@@ -21,6 +22,11 @@ export const KUBERNETES_CONFIG_LIST: Config[] = [{
name
:
'kubernetes.container.image'
,
name
:
'kubernetes.container.image'
,
lable
:
'kubernetes.container.image'
,
lable
:
'kubernetes.container.image'
,
placeholder
:
'dlink'
,
placeholder
:
'dlink'
,
},{
name
:
'kubernetes.rest-service.exposed.type'
,
lable
:
'kubernetes.rest-service.exposed.type'
,
placeholder
:
'NodePort'
,
defaultValue
:
'NodePort'
,
}];
}];
export
const
FLINK_CONFIG_LIST
:
Config
[]
=
[{
export
const
FLINK_CONFIG_LIST
:
Config
[]
=
[{
name
:
'jobmanager.memory.process.size'
,
name
:
'jobmanager.memory.process.size'
,
...
...
dlink-web/src/pages/Welcome.tsx
View file @
be150742
...
@@ -707,6 +707,9 @@ export default (): React.ReactNode => {
...
@@ -707,6 +707,9 @@ export default (): React.ReactNode => {
<
li
>
<
li
>
<
Link
>
新增 实时自动告警
</
Link
>
<
Link
>
新增 实时自动告警
</
Link
>
</
li
>
</
li
>
<
li
>
<
Link
>
优化 K8S Application 提交配置
</
Link
>
</
li
>
</
ul
>
</
ul
>
</
Paragraph
>
</
Paragraph
>
</
Timeline
.
Item
>
</
Timeline
.
Item
>
...
...
Write
Preview
Markdown
is supported
0%
Try again
or
attach a new file
Attach a file
Cancel
You are about to add
0
people
to the discussion. Proceed with caution.
Finish editing this message first!
Cancel
Please
register
or
sign in
to comment