ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Apache Pulsar Functions 调试实战指南:从单元测试到 localrun 与 CLI 全流程

2026/9/23 6:22:45 拓冰建站 浏览量
Apache Pulsar Functions 调试实战指南:从单元测试到 localrun 与 CLI 全流程 Apache Pulsar Functions 调试实战指南从单元测试到 localrun 与 CLI 全流程【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsar本篇指南以 Apache Pulsar 官方文档 functions-debug 为核心骨架系统讲解调试 Pulsar Functions 的五种方法捕获 stderr 日志、单元测试、localrun 本地运行模式、日志主题log topic以及 Functions CLI 子命令。读完本文你将掌握从“函数无法启动”的快速定位到“在 IDE 中打断点调试真实数据”再到“线上运行时状态与指标排查”的完整调试技能。调试方法总览Apache Pulsar Functions 提供了五种互补的调试手段覆盖了函数生命周期中从开发、启动到运行监控的各个阶段Captured stderr查看函数启动信息与捕获的 stderr 输出用于定位启动失败原因。单元测试像测试普通函数一样测试 Pulsar Function快速验证处理逻辑。localrun 模式在本地以线程方式运行函数实例可在 IDE 中打断点、用真实数据逐步调试。日志主题把函数内定义的日志写入指定 Pulsar topic用消费者检查运行期日志。Functions CLI通过pulsar-admin functions的get、status、stats、list、trigger子命令在线排查。Captured stderr定位函数启动失败函数启动信息与被捕获的 stderr 输出会写入以下路径logs/functions/tenant/namespace/function/function-instance.log即日志文件按租户/命名空间/函数名/函数名-实例ID的目录结构组织。当函数无法启动时先查看对应实例的这份日志通常能直接找到类加载异常、依赖缺失或配置错误等根因。例如一个属于public/default命名空间、名为my-function且实例 ID 为 0 的函数其日志位于logs/functions/public/default/my-function/my-function-0.log使用单元测试调试函数Pulsar Function 本质上是一个“有输入、有输出”的函数因此你可以像测试任何普通函数一样对它进行单元测试。项目中的实际测试代码也印证了这一模式——例如 pulsar-functions/instance 模块内大量的*Test.java测试文件以及整个 Pulsar 仓库约定使用 TestNG 作为测试框架。测试原生 Java Function先看一个使用 JDK 自带java.util.function.Function实现的函数import java.util.function.Function; public class JavaNativeExclamationFunction implements FunctionString, String { Override public String apply(String input) { return String.format(%s!, input); } }对应的单元测试极其简单直接调用apply断言输出Test public void testJavaNativeExclamationFunction() { JavaNativeExclamationFunction exclamation new JavaNativeExclamationFunction(); String output exclamation.apply(foo); Assert.assertEquals(output, foo!); }测试 Pulsar API 的 Function 接口更多时候函数实现的是org.apache.pulsar.functions.api.Function接口。该接口定义在 pulsar-functions/api-java/src/main/java/org/apache/pulsar/functions/api/Function.java核心方法为O process(I input, Context context)此外还提供了initialize(Context)与close()两个默认为空的生命周期钩子方法分别在实例启动与停止时被调用一次。import org.apache.pulsar.functions.api.Context; import org.apache.pulsar.functions.api.Function; public class ExclamationFunction implements FunctionString, String { Override public String process(String input, Context context) { return String.format(%s!, input); } }由于process多了一个Context参数测试时需要 Mock 掉它。Pulsar 使用 TestNG因此可以直接借助org.mockito.Mockito.mockTest public void testExclamationFunction() { ExclamationFunction exclamation new ExclamationFunction(); String output exclamation.process(foo, mock(Context.class)); Assert.assertEquals(output, foo!); }Context接口定义在 pulsar-functions/api-java/src/main/java/org/apache/pulsar/functions/api/Context.java它向函数暴露输入主题、输出主题、用户配置、消息 ID 等上下文信息在单元测试中用 Mock 隔离外部依赖是最佳实践。使用 localrun 模式调试什么是 localrun 模式当以 localrun 模式运行 Pulsar Function 时函数实例会在你本地机器上以**线程thread**的形式被拉起。该模式下函数会向一个真实的 Pulsar 集群消费、生产真实数据行为与在集群中运行时完全一致因此特别适合在 IDE 里设置断点、单步跟踪用真实消息验证处理逻辑。注意localrun 调试目前仅支持 Java 编写的 Pulsar Function且需要 Pulsar 2.4.0 或更高版本。虽然在 2.4.0 之前的版本中 localrun 已存在但无法通过编程方式调试或作为线程运行函数。从源码实现看localrun 的核心类是 pulsar-functions/localrun/src/main/java/org/apache/pulsar/functions/LocalRunner.java。它内部维护一个RuntimeSpawner列表start(boolean blocking)方法根据配置选择两种运行方式Thread 模式ThreadRuntimeFactory默认方式Java 函数以线程在本地 JVM 中运行方便 IDE 断点调试。Process 模式ProcessRuntimeFactory以独立进程方式运行。默认的 broker 服务地址为pulsar://localhost:6650Web 服务地址为http://localhost:8080与本地 standalone 集群默认端口一致。编程方式启动函数通过LocalRunner.builder()构建并启动FunctionConfig functionConfig new FunctionConfig(); functionConfig.setName(functionName); functionConfig.setInputs(Collections.singleton(sourceTopic)); functionConfig.setClassName(ExclamationFunction.class.getName()); functionConfig.setRuntime(FunctionConfig.Runtime.JAVA); functionConfig.setOutput(sinkTopic); LocalRunner localRunner LocalRunner.builder().functionConfig(functionConfig).build(); localRunner.start(true);start(true)表示阻塞运行线程会阻塞等待start(false)表示非阻塞返回。以下是一个包含main方法的完整示例可直接在 IDE 中运行并在process方法内打断点public class ExclamationFunction implements FunctionString, String { Override public String process(String s, Context context) throws Exception { return s !; } public static void main(String[] args) throws Exception { FunctionConfig functionConfig new FunctionConfig(); functionConfig.setName(exclamation); functionConfig.setInputs(Collections.singleton(input)); functionConfig.setClassName(ExclamationFunction.class.getName()); functionConfig.setRuntime(FunctionConfig.Runtime.JAVA); functionConfig.setOutput(output); LocalRunner localRunner LocalRunner.builder().functionConfig(functionConfig).build(); localRunner.start(false); }提示从LocalRunner源码看它不仅支持FunctionConfig还支持SourceConfig与SinkConfig也就是说 localrun 同样可以用于调试 Source 和 Sink 连接器。此外通过 builder 还可以设置brokerServiceUrl、webServiceUrl、metricsPortStartPrometheus 指标端口、runtimeEnvTHREAD/PROCESS等参数其中默认情况下 Java 函数按线程模式运行。添加 localrun 依赖要在代码中编程使用 localrun需要添加如下 Maven 依赖dependency groupIdorg.apache.pulsar/groupId artifactIdpulsar-functions-local-runner/artifactId version${pulsar.version}/version /dependency本地调试的运行流程建议先用bin/pulsar standalone启动一个 standalone 集群或连接已有集群然后在 IDE 中把process方法设断点向输入 topic 生产消息观察函数对真实数据的逐条处理。此外LocalRunner还支持通过--runtime指定 Process 模式RuntimeEnv.PROCESS此时会在本地以独立进程方式模拟集群中函数实例的运行形态。使用日志主题Log Topic调试原理与编程模型在 Pulsar Functions 中函数内定义的日志信息可以发送到指定的 log topic。你可以配置消费者消费该 topic 中的消息来查看日志。从源码层面看日志发送由 pulsar-functions/instance/src/main/java/org/apache/pulsar/functions/instance/LogAppender.java 实现它实现 Log4j2 的Appender接口将日志事件格式化为消息后通过Producer异步发送到 log topic。生产端默认启用 LZ4 压缩与批量发送batchingMaxPublishDelay为 100ms并给每条消息写入function属性。在函数中写日志以下LoggingFunction通过context.getLogger()获取 SLF4J Logger并输出带消息 ID 的日志import org.apache.pulsar.functions.api.Context; import org.apache.pulsar.functions.api.Function; import org.slf4j.Logger; public class LoggingFunction implements FunctionString, Void { Override public void apply(String input, Context context) { Logger LOG context.getLogger(); String messageId new String(context.getMessageId()); if (input.contains(danger)) { LOG.warn(A warning was received in message {}, messageId); } else { LOG.info(Message {} received\nContent: {}, messageId, input); } return null; } }getLogger()在运行实例的 ContextImpl 中实现返回与该函数实例绑定的 Logger当配置了 log topic 时该 Logger 会将日志路由到LogAppender并写入 topic。指定 log topic创建函数时通过--log-topic指定日志要发往的主题$ bin/pulsar-admin functions create \ --log-topic persistent://public/default/logging-function-logs \ # Other function configs日志消息的属性发往 log topic 的消息包含以下属性便于在消费端进行过滤与定位loglevel—— 日志消息的级别如 INFO、WARN、ERROR。fqn—— 产生该日志消息的函数全限定名Fully Qualified Function Name。instance—— 产生该日志消息的函数实例 ID。这三个属性与LogAppender中的常量一一对应LOG_LEVEL、FQN、INSTANCE在发送时被写入消息属性。你可以用pulsar-client consume或任意 Pulsar 消费者订阅该 topic按loglevel过滤 ERROR/WARN或按fqn、instance定位到具体函数与实例从而在函数运行期持续观测日志输出。使用 Functions CLI 调试借助 Pulsar Functions CLI可以使用以下子命令调试 Pulsar Functionsgetstatusstatslisttrigger这些子命令的实现位于 pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdFunctions.java其背后的管理接口则定义在 pulsar-client-admin-api 的 Functions 管理 API 中。get查看函数配置信息获取某个 Pulsar Function 的完整配置信息。$ pulsar-admin functions get optionsOptionsFlag描述--fqfnPulsar Function 的全限定函数名FQFN。--namePulsar Function 的名称。--namespacePulsar Function 所在的命名空间。--tenantPulsar Function 所在的租户。提示--fqfn由--name、--namespace和--tenant组合而成因此你可以二选一要么指定--fqfn要么分别指定--name、--namespace和--tenant。示例通过--fqfn获取函数信息$ ./bin/pulsar-admin functions get public/default/ExclamationFunctio6或者分别指定--name、--namespace与--tenant$ ./bin/pulsar-admin functions get \ --tenant public \ --namespace default \ --name ExclamationFunctio6输出示例展示了输入、输出、运行时等配置{ tenant: public, namespace: default, name: ExclamationFunctio6, className: org.example.test.ExclamationFunction, inputSpecs: { persistent://public/default/my-topic-1: { isRegexPattern: false } }, output: persistent://public/default/test-1, processingGuarantees: ATLEAST_ONCE, retainOrdering: false, userConfig: {}, runtime: JAVA, autoAck: true, parallelism: 1 }status查看函数运行状态检查某个 Pulsar Function 的当前运行状态。$ pulsar-admin functions status optionsOptionsFlag描述--fqfnPulsar Function 的全限定函数名FQFN。--instance-idPulsar Function 的实例 ID。若未指定--instance-id则返回所有实例的状态。--namePulsar Function 的名称。--namespacePulsar Function 所在的命名空间。--tenantPulsar Function 所在的租户。示例$ ./bin/pulsar-admin functions status \ --tenant public \ --namespace default \ --name ExclamationFunctio6 \输出示例展示实例数、运行实例数、已接收消息、成功处理消息、系统异常、平均延迟等{ numInstances : 1, numRunning : 1, instances : [ { instanceId : 0, status : { running : true, error : , numRestarts : 0, numReceived : 1, numSuccessfullyProcessed : 1, numUserExceptions : 0, latestUserExceptions : [ ], numSystemExceptions : 0, latestSystemExceptions : [ ], averageLatency : 0.8385, lastInvocationTime : 1557734137987, workerId : c-standalone-fw-23ccc88ef29b-8080 } } ] }重点关注字段running表示实例是否在运行numRestarts反映实例重启次数启动失败排查的关键指标numUserExceptions与numSystemExceptions分别统计用户代码异常和系统级异常对应的latestUserExceptions、latestSystemExceptions则给出最近一次异常的详细信息。stats查看函数运行指标获取某个 Pulsar Function 当前的统计指标。$ pulsar-admin functions stats optionsOptionsFlag描述--fqfnPulsar Function 的全限定函数名FQFN。--instance-idPulsar Function 的实例 ID。若未指定--instance-id则返回所有实例的指标。--namePulsar Function 的名称。--namespacePulsar Function 所在的命名空间。--tenantPulsar Function 所在的租户。示例$ ./bin/pulsar-admin functions stats \ --tenant public \ --namespace default \ --name ExclamationFunctio6 \输出示例{ receivedTotal : 1, processedSuccessfullyTotal : 1, systemExceptionsTotal : 0, userExceptionsTotal : 0, avgProcessLatency : 0.8385, 1min : { receivedTotal : 0, processedSuccessfullyTotal : 0, systemExceptionsTotal : 0, userExceptionsTotal : 0, avgProcessLatency : null }, lastInvocation : 1557734137987, instances : [ { instanceId : 0, metrics : { receivedTotal : 1, processedSuccessfullyTotal : 1, systemExceptionsTotal : 0, userExceptionsTotal : 0, avgProcessLatency : 0.8385, 1min : { receivedTotal : 0, processedSuccessfullyTotal : 0, systemExceptionsTotal : 0, userExceptionsTotal : 0, avgProcessLatency : null }, lastInvocation : 1557734137987, userMetrics : { } } } ] }stats与status的区别在于stats侧重吞吐量与延迟等指标如receivedTotal、avgProcessLatency并含最近 1 分钟窗口数据status侧重运行状态与异常计数。两者结合即可全面评估函数实例的健康度。list列出命名空间下的函数列出某个特定租户与命名空间下运行的所有 Pulsar Functions。$ pulsar-admin functions list optionsOptionsFlag描述--namespacePulsar Function 所在的命名空间。--tenantPulsar Function 所在的租户。示例$ ./bin/pulsar-admin functions list \ --tenant public \ --namespace default输出示例public租户、default命名空间下运行着三个函数ExclamationFunctio1 ExclamationFunctio2 ExclamationFunctio3trigger触发函数并验证使用指定值触发某个 Pulsar Function。该命令会模拟函数的执行过程并验证结果。$ pulsar-admin functions trigger optionsOptionsFlag描述--fqfnPulsar Function 的全限定函数名FQFN。--namePulsar Function 的名称。--namespacePulsar Function 所在的命名空间。--tenantPulsar Function 所在的租户。--topicPulsar Function 消费的 topic 名称。--trigger-file包含触发数据的文件路径。--trigger-value用于触发 Pulsar Function 的值。示例$ ./bin/pulsar-admin functions trigger \ --tenant public \ --namespace default \ --name ExclamationFunctio6 \ --topic persistent://public/default/my-topic-1 \ --trigger-value hello pulsar functions输出示例This is my function!注意使用--topic选项时必须指定完整的 topic 名称即带persistent://tenant/namespace/topic前缀。否则会出现如下错误Function in trigger function has unidentified topic Reason: Function in trigger function has unidentified topic调试方法选择建议调试阶段推荐方法适用场景逻辑验证单元测试无需外部依赖快速验证process处理逻辑与边界情况启动失败排查Captured stderr函数无法启动、类加载异常、依赖缺失真实数据联调localrun 模式在 IDE 中断点调试使用真实消息流验证运行期观测Log topic分布式环境下的远程日志收集与分析在线状态排查Functions CLI查询函数配置、状态、指标手动触发验证五种方法覆盖了从开发期到生产期的完整调试链路开发阶段先用单元测试锁定逻辑正确性再通过 localrun 结合真实数据做集成级调试函数部署上线后用 Captured stderr 与status/stats子命令监控健康度用 log topic 收集运行日志最后用trigger对线上函数做快速冒烟验证。相关示例源码可在 pulsar-functions/java-examples 与 pulsar-functions/localrun 中继续深入阅读。【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址: https://gitcode.com/gh_mirrors/pulsar28/pulsar创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考