Java开发

最近更新时间: 2026-06-30 15:06:00

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.hostTCT服务端地址,可以在 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