JDK 版本
推荐使用 JDK 1.8 或者以上版本,对于高于 1.8 的 JDK 版本,需要额外添加以下依赖包:
<dependency>
<groupId>javax.xml.bind</groupId>
<artifactId>jaxb-api</artifactId>
<version>2.3.0</version>
</dependency>
添加依赖
找到 Maven 所使用的配置文件 settings.xml,一般为 ~/.m2/settings.xml,添加 TCT Maven 地址。
<?xml version="1.0" encoding="UTF-8"?>
<settings xmlns="http://maven.apache.org/SETTINGS/1.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/SETTINGS/1.0.0 http://maven.apache.org/xsd/settings-1.0.0.xsd">
<pluginGroups></pluginGroups>
<proxies></proxies>
<servers></servers>
<mirrors></mirrors>
<profiles>
<profile>
<id>nexus</id>
<repositories>
<repository>
<id>central</id>
<url>http://repo1.maven.org/maven2</url>
<releases>
<enabled>true</enabled>
</releases>
<snapshots>
<enabled>true</enabled>
</snapshots>
</repository>
</repositories>
<pluginRepositories>
<pluginRepository>
<id>central</id>
<url>http://repo1.maven.org/maven2</url>
<releases>
<enabled>true</enabled>
</releases>
<snapshots>
<enabled>true</enabled>
</snapshots>
</pluginRepository>
</pluginRepositories>
</profile>
<profile>
<id>tct</id>
<repositories>
<repository>
<id>tct</id>
<name>tct</name>
<url>https://mirrors.cloud.tencent.com/nexus/repository/maven-public/</url>
<releases>
<enabled>true</enabled>
</releases>
<snapshots>
<enabled>true</enabled>
</snapshots>
</repository>
</repositories>
</profile>
</profiles>
<activeProfiles>
<activeProfile>nexus</activeProfile>
<activeProfile>tct</activeProfile>
</activeProfiles>
</settings>
然后在 pom.xml 文件中添加 TCT spring boot starter 依赖。
<dependency>
<groupId>com.tencent.cloud </groupId>
<artifactId>tct-spring-boot-starter </artifactId>
<version>2.1.0-rc8</version>
</dependency>
任务开发
简单任务
编写 TCT 任务,只需要实现 TCT 提供的 com.tencent.cloud.task.sdk.client.spi.ExecutableTask 接口,在 execute 方法中实现任务执行逻辑,SDK 内部通过反射机制,生成任务对象实例,并执行 execute 方法。如下所示, SleepTask 是一个 sleep 10 秒的简单任务。
注意,将任务申明为 Bean 才能被 SDK 自动发现并作为预置任务上报给 TCT,如果不申明为 Bean 任务也可以使用但是在 TCT 控制台上创建任务时需要手动输入完整的类名。
import com.tencent.cloud.task.sdk.client.LogReporter;
import com.tencent.cloud.task.sdk.client.model.ExecutableTaskData;
import com.tencent.cloud.task.sdk.client.model.ProcessResult;
import com.tencent.cloud.task.sdk.client.spi.ExecutableTask;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;
import java.lang.invoke.MethodHandles;
@Component
public class SleepTask implements ExecutableTask {
private static final Logger LOG = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass());
@Override
public ProcessResult execute(ExecutableTaskData taskData) {
LOG.info("Run sleep task, taskMeta: {}", taskData.getTaskMeta().toString());
try {
Thread.sleep(10 * 1000L);
return ProcessResult.newSuccessResult();
} catch (InterruptedException e) {
LogReporter.log(taskData, "Task is terminated.");
return ProcessResult.newCancelledResult();
} catch (Throwable e) {
LogReporter.log(taskData, String.format("Exception when sleep: %s", e.getMessage()));
return ProcessResult.newFailResult();
}
}
}
可停止任务
在 TCT 控制台中,我们可以停止一个执行中的任务,为了使任务可停止,编写任务逻辑的时候,需要实现 com.tencent.cloud.task.sdk.client.spi.TerminableTask 接口,在 cancel 方法中实现任务的停止逻辑,并返回停止结果。
如下所示,SleepTask 中我们通过 Future 的 cancel 方法停止 sleep 任务,这时候任务执行线程会收到中断信号,抛出中断异常 InterruptedException,示例代码中捕获了 InterruptedException 异常并返回任务终止成功。
import com.tencent.cloud.task.sdk.client.LogReporter;
import com.tencent.cloud.task.sdk.client.model.ExecutableTaskData;
import com.tencent.cloud.task.sdk.client.model.ProcessResult;
import com.tencent.cloud.task.sdk.client.model.TerminateResult;
import com.tencent.cloud.task.sdk.client.remoting.TaskExecuteFuture;
import com.tencent.cloud.task.sdk.client.spi.ExecutableTask;
import com.tencent.cloud.task.sdk.client.spi.TerminableTask;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;
import java.lang.invoke.MethodHandles;
@Component
public class SleepTask implements ExecutableTask, TerminableTask {
private static final Logger LOG = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass());
@Override
public ProcessResult execute(ExecutableTaskData taskData) {
LOG.info("Run sleep task, taskMeta: {}", taskData.getTaskMeta().toString());
try {
Thread.sleep(10 * 1000L);
return ProcessResult.newSuccessResult();
} catch (InterruptedException e) {
LogReporter.log(taskData, "Task is terminated.");
return ProcessResult.newCancelledResult();
} catch (Throwable e) {
LogReporter.log(taskData, String.format("Exception when sleep: %s", e.getMessage()));
return ProcessResult.newFailResult();
}
}
@Override
public TerminateResult cancel(TaskExecuteFuture taskExecuteFuture, ExecutableTaskData executableTaskData) {
taskExecuteFuture.cancel(true);
return TerminateResult.newTerminateSuccessResult();
}
}
TCT SDK 扫描预置任务的时候会通过判断任务类是否实现了 TerminableTask 接口判断任务是否支持停止,TCT 控制台中只有支持停止的任务允许执行停止操作。
任务停止原理
首先,我们需要了解 Java 体系内,如何停止一个执行中的任务。向一个 Alive 状态的 Thread 发送中断信号,不一定能中断线程的执行。中断信号只是向运行的线程一个建议,告诉它有外界希望中断它,至于线程接受到信号后要做出何种反应,完全由线程及运行状态自身决定。
Java 提供的 API 中,对于中断信号,通常存在两种响应形态:
- 设置中断状态标志位,通过 Thread.isInterrupted() 进行判断。
- 当执行任务的线程处于 BLOCKED 状态(例如调用了 wait、sleep、join 等方法)时,向线程发送中断信号,通常会抛出中断异常,如 java.lang.InterruptedException。Java 中常见的中断异常有 java.lang.InterruptedException、 java.io.InterruptedIOException、java.nio.channels.ClosedByInterruptException。
TCT 对用户侧提供了中断任务执行线程的 API 来向执行线程发送中断信号:
cancel(TaskExecuteFuture future, ExecutableTaskData tasData)
此方法中暴露 Future 对象,底层是对当前任务提交执行后返回的java.util.concurrent.future的封装。 通过调用 future.cancel(boolean) 向执行任务的线程发送中断信号。 当在控制台操作停止任务时,cancel 方法将会被调用。
因此,我们在实现任务执行逻辑的时候,需要判断中断标志位和捕获中断异常来实现可停止的任务,如下所示,加入我们的任务需要循环处理一批数据,在每一次循环的时候我们都判断:
@Override
public ProcessResult execute(ExecutableTaskData taskData) {
try {
List<String> dataset = Arrays.asList("id1", "id2", "id3");
for(String data : dataset) {
if (Thread.currentThread().isInterrupted()) {
return ProcessResult.newCancelledResult();
}
// 数据处理逻辑。。。
}
} catch (InterruptedException e) {
LogReporter.log(taskData, "Task is terminated.");
return ProcessResult.newCancelledResult();
} catch (Throwable e) {
LogReporter.log(taskData, String.format("Exception when sleep: %s", e.getMessage()));
return ProcessResult.newFailResult();
}
}
除了通过中断信号实现可停止任务外,我们也可以通过 execute 和 cancel 两个接口配合实现业务的逻辑终止,例如通过一个 isCancelled 字段标识业务逻辑是否被终止,在 cancel 方法中设置它,在 execute 中执行业务逻辑时检查它。
import com.tencent.cloud.task.sdk.client.model.ExecutableTaskData;
import com.tencent.cloud.task.sdk.client.model.ProcessResult;
import com.tencent.cloud.task.sdk.client.model.TerminateResult;
import com.tencent.cloud.task.sdk.client.remoting.TaskExecuteFuture;
import com.tencent.cloud.task.sdk.client.spi.ExecutableTask;
import com.tencent.cloud.task.sdk.client.spi.TerminableTask;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.stereotype.Component;
import java.lang.invoke.MethodHandles;
import java.util.Arrays;
import java.util.List;
import java.util.concurrent.atomic.AtomicBoolean;
@Component
public class SleepTask implements ExecutableTask, TerminableTask {
private static final Logger LOG = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass());
private final AtomicBoolean isCancelled = new AtomicBoolean(false);
@Override
public ProcessResult execute(ExecutableTaskData taskData) {
List<String> dataset = Arrays.asList("id1", "id2", "id3");
for (String data : dataset) {
if (isCancelled.get()) {
return ProcessResult.newCancelledResult();
}
// 数据处理逻辑。。。
}
return ProcessResult.newSuccessResult();
}
@Override
public TerminateResult cancel(TaskExecuteFuture taskExecuteFuture, ExecutableTaskData executableTaskData) {
// 设置终止状态,终止成功
isCancelled.set(true);
// 返回终止成功
return TerminateResult.newTerminateSuccessResult();
}
}
任务工厂
任务工厂是 TCT 用来生成任务实例的工厂类,TCT 提供了默认的工厂类 com.tencent.cloud.task.sdk.client.DefaultTaskFactory,也支持用户自定义任务工厂类。
默认工厂
默认工厂 DefaultTaskFactory 通过 Java 的反射机制来生成任务对象的实例,如下代码所示:
public class DefaultTaskFactory implements ExecutableTaskFactory {
private static final Logger LOG = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass());
private final ClassLoader classLoader;
public DefaultTaskFactory(ClassLoader classLoader) {
this.classLoader = classLoader;
}
@Override
public ExecutableTask newExecutableTask(ExecutableTaskData taskData) throws InstancingException {
String taskName = taskData.getTaskContent();
if (LOG.isDebugEnabled()) {
LOG.debug("producing instance of ExecutableTask: '" + taskName + "'");
}
try {
Class<?> taskClass = Class.forName(taskName, true, classLoader);
if (!ExecutableTask.class.isAssignableFrom(taskClass)) {
throw new InstancingException("Problem instancing ExecutableTask, "
+ "Caused by task Class name '" + ExecutableTask.class.getName()
+ "' is not AssignableFrom Class '" + taskName + "'");
}
return (ExecutableTask) taskClass.newInstance();
} catch (ClassNotFoundException t) {
throw new InstancingException("Class '" + taskName + "' is not found", t);
} catch (Exception e) {
if (e instanceof InstancingException) {
throw (InstancingException) e;
}
throw new InstancingException("Problem instancing ExecutableTask,"
+ " task Class is '" + taskName + "'", e);
}
}
}
自定义任务工厂
普通 Java 工厂
创建自定义任务工厂,只需要实现 com.tencent.cloud.task.sdk.client.spi.ExecutableTaskFactory 接口,如下示例所示,我们创建了 SimpleExecuteTaskFactory 类,它扩展了默认的 DefaultTaskFactory 工厂类,DefaultTaskFactory 实现了 ExecutableTaskFactory 接口。
package com.tencent.cloud.task.factory
public class SimpleExecuteTaskFactory extends DefaultTaskFactory {
private static final Logger LOG = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass());
public SimpleExecuteTaskFactory() {
super(Thread.currentThread().getContextClassLoader());
}
public SimpleExecuteTaskFactory(ClassLoader classLoader) {
super(classLoader);
}
@Override
public ExecutableTask newExecutableTask(ExecutableTaskData taskData) throws InstancingException {
LOG.info("generate task: {}", taskData.getTaskContent());
return super.newExecutableTask(taskData);
}
}
创建好工厂类后,我们需要修改任务应用(即执行器)的启动配置,在 application.yml 里配置任务工厂类:
tct:
client:
properties:
"task.factory.name": "com.tencent.cloud.task.factory.SimpleExecuteTask
Spring 框架工厂
在 Spring 框架中,我们可以将任务注册成 bean,然后自定义任务工厂从应用上下文中获取这些任务 bean。
@Component
public class SpringExecuteTaskFactory implements ExecutableTaskFactory, ApplicationContextAware {
private static final Logger LOG = LoggerFactory.getLogger(MethodHandles.lookup().lookupClass());
private ApplicationContext applicationContext;
private final ExecutableTaskFactory defaultFactory = new DefaultTaskFactory(Thread.currentThread().getContextClassLoader());
@Override
public ExecutableTask newExecutableTask(ExecutableTaskData executableTaskData) throws InstancingException {
try {
ExecutableTask executableTask = (ExecutableTask)applicationContext.getBean(Class.forName(executableTaskData.getTaskContent()));
LOG.info("generate executableTask bean SpringExecutableTaskFactory. taskName: {}", executableTaskData.getTaskContent());
return executableTask;
} catch (Throwable t) {
return defaultFactory.newExecutableTask(executableTaskData);
}
}
@Override
public void setApplicationContext(ApplicationContext applicationContext) throws BeansException {
this.applicationContext = applicationContext;
}
同样需要修改配置指定工厂类:
tct:
client:
properties:
"task.factory.name": "com.tencent.cloud.task.factory.SpringExecuteTaskFactory"
任务配置
TCT 任务应用(即执行器)提供以下配置项供用户配置:
| 配置项 | 说明 |
| tct.enabled | 是否开启 TCT 任务调度,只有该配置项为 true 才会启用 TCT 功能。 |
| tct.server.host | TCT服务端地址,可以在 TCT 部署组详情中获取。 |
| tct.server.port | |
| tct.client.groupId | 部署组 ID,如果任务应用通过 TSF 部署,该配置可以不填,TSF 会自动注入部署组 ID。其他情况下需要手动填充在 TCT 控制台创建好的部署组的 ID,例如默认的部署组填 default。 |
| tct.client.instanceId | 任务应用实例(执行器实例)的 ID,如果任务应用通过 TSF 部署,在部署时会自动注入,该配置可不填。其他情况下需要自行配置,注意这里的 ID 当前需要在部署组范围内唯一。 |
| tct.client.accessKey | 用于认证和鉴权,在TCS 控制台获取,当前登录用户 - 账号信息 - API密钥管理。 |
| tct.client.secretKey | |
| tct.client.environments | 指定该任务应用(执行器)具备怎样的执行环境,即支持执行什么类型的任务,例如 Java、Python、Shell、External。 |
可以通过以下几种方式配置这些配置项,并且几种方式的优先级如下,高优先级的配置会覆盖低优先级的配置:环境变量 > 命令行参数 > application.yml。
application.yml
tct:
enabled: true
server:
host: server.chongqing.tct
port: 28000
client:
groupId: BjFnVcXkVw
instanceId: tct-demo-ins1
accessKey: xxx
secretKey: xxx
environments:
- Java
- Shell
- Python
命令行参数
java \
-Dtct.server.host=10.0.8.24 \
-Dtct.server.port=28000 \
-Dtct.client.groupId=BjFnVcXkVw \
-Dtct.client.instanceId=tct-demo-ins1 \
-Dtct.client.accessKey=xxx \
-Dtct.client.secretKey=xxx \
-jar tct-demo.jar
环境变量
可以通过其他方式注入环境变量(例如 k8s Pod 中指定 env),也可以通过启动参数注定环境变量,例如:
java \
-Dtct_server_host=10.0.8.24 \
-Dtct_server_port=28000 \
-Dtct_group_id=BjFnVcXkVw \
-Dtct_instance_id=tct-demo-ins1 \
-Dtct_access_key=xxx \
-Dtct_secret_key=xxx \
-Dtct_environments=Java,Python,Shell \
-jar tct-demo.jar