- xml - AJAX/Jquery XML 解析
- 具有多重继承的 XML 模式
- .net - 枚举序列化 Json 与 XML
- XML 简单类型、简单内容、复杂类型、复杂内容
我刚开始学习 channel 。我正在使用汇合的 kafka 消费者来创建功能性消费者。我想要完成的是将消息发送到缓冲 channel (2,000)...然后使用管道将 channel 中的消息写入 redis。我已经通过执行 println
来让消费者部分工作了一条一条地发送消息,直到它到达偏移量的末尾,但是当我尝试添加一个 channel 时,它似乎命中了 default:
switch
中的案例然后就卡住了。
我似乎也没有正确填写 channel ?这fmt.Println("count is: ", len(redisChnl))
总是打印 0
这是我目前所拥有的:
// Example function-based high-level Apache Kafka consumer
package main
import (
"fmt"
"github.com/confluentinc/confluent-kafka-go/kafka"
"os"
"os/signal"
"syscall"
"time"
"encoding/json"
"regexp"
"github.com/go-redis/redis"
"encoding/binary"
)
var client *redis.Client
func init() {
client = redis.NewClient(&redis.Options{
Addr: ":6379",
DialTimeout: 10 * time.Second,
ReadTimeout: 30 * time.Second,
WriteTimeout: 30 * time.Second,
PoolSize: 10,
PoolTimeout: 30 * time.Second,
})
client.FlushDB()
}
type MessageFormat struct {
MetricValueNumber float64 `json:"metric_value_number"`
Path string `json:"path"`
Cluster string `json:"cluster"`
Timestamp time.Time `json:"@timestamp"`
Version string `json:"@version"`
Host string `json:"host"`
MetricPath string `json:"metric_path"`
Type string `json:"string"`
Region string `json:"region"`
}
//func redis_pipeline(ky string, vl string) {
// pipe := client.Pipeline()
//
// exec := pipe.Set(ky, vl, time.Hour)
//
// incr := pipe.Incr("pipeline_counter")
// pipe.Expire("pipeline_counter", time.Hour)
//
// // Execute
// //
// // INCR pipeline_counter
// // EXPIRE pipeline_counts 3600
// //
// // using one client-server roundtrip.
// _, err := pipe.Exec()
// fmt.Println(incr.Val(), err)
// // Output: 1 <nil>
//}
func main() {
sigchan := make(chan os.Signal, 1)
signal.Notify(sigchan, syscall.SIGINT, syscall.SIGTERM)
c, err := kafka.NewConsumer(&kafka.ConfigMap{
"bootstrap.servers": "kafka.com:9093",
"group.id": "testehb",
"security.protocol": "ssl",
"ssl.key.location": "/Users/key.key",
"ssl.certificate.location": "/Users/cert.cert",
"ssl.ca.location": "/Users/ca.pem",
})
if err != nil {
fmt.Fprintf(os.Stderr, "Failed to create consumer: %s\n", err)
os.Exit(1)
}
fmt.Printf("Created Consumer %v\n", c)
err = c.SubscribeTopics([]string{"jmx"}, nil)
redisMap := make(map[string]string)
redisChnl := make(chan []byte, 2000)
run := true
for run == true {
select {
case sig := <-sigchan:
fmt.Printf("Caught signal %v: terminating\n", sig)
run = false
default:
ev := c.Poll(100)
if ev == nil {
continue
}
switch e := ev.(type) {
case *kafka.Message:
//fmt.Printf("%% Message on %s:\n%s\n",
// e.TopicPartition, string(e.Value))
if e.Headers != nil {
fmt.Printf("%% Headers: %v\n", e.Headers)
}
str := e.Value
res := MessageFormat{}
json.Unmarshal([]byte(str), &res)
fmt.Println("size", binary.Size([]byte(str)))
host:= regexp.MustCompile(`^([^.]+)`).FindString(res.MetricPath)
redisMap[host] = string(str)
fmt.Println("count is: ", len(redisChnl)) //this always prints "count is: 0"
redisChnl <- e.Value //I think this is the write way to put the messages in the channel?
case kafka.PartitionEOF:
fmt.Printf("%% Reached %v\n", e)
case kafka.Error:
fmt.Fprintf(os.Stderr, "%% Error: %v\n", e)
run = false
default:
fmt.Printf("Ignored %v\n", e)
}
<- redisChnl // I thought I could just empty the channel like this once the buffer is full?
}
}
fmt.Printf("Closing consumer\n")
c.Close()
}
--------编辑--------
好的,我想我是通过移动 <- redisChnl
让它工作的里面default
, 但现在我看到 count before read
和 count after read
在default
里面总是打印 2,000
...我本以为第一个count before read = 2,000
然后 count after read = 0
因为那时 channel 将是空的??
select {
case sig := <-sigchan:
fmt.Printf("Caught signal %v: terminating\n", sig)
run = false
default:
ev := c.Poll(100)
if ev == nil {
continue
}
switch e := ev.(type) {
case *kafka.Message:
//fmt.Printf("%% Message on %s:\n%s\n",
// e.TopicPartition, string(e.Value))
if e.Headers != nil {
fmt.Printf("%% Headers: %v\n", e.Headers)
}
str := e.Value
res := MessageFormat{}
json.Unmarshal([]byte(str), &res)
//fmt.Println("size", binary.Size([]byte(str)))
host:= regexp.MustCompile(`^([^.]+)`).FindString(res.MetricPath)
redisMap[host] = string(str)
go func() {
redisChnl <- e.Value
}()
case kafka.PartitionEOF:
fmt.Printf("%% Reached %v\n", e)
case kafka.Error:
fmt.Fprintf(os.Stderr, "%% Error: %v\n", e)
run = false
default:
fmt.Println("count before read: ", len(redisChnl))
fmt.Printf("Ignored %v\n", e)
<-redisChnl
fmt.Println("count after read: ", len(redisChnl)) //would've expected this to be 0
}
}
最佳答案
我认为简化此代码的更大方法是将管道分成多个 goroutine。
channel 的优点是多人可以同时在上面书写和阅读。在这个例子中,这可能意味着有一个 go 例程在 channel 上排队,另一个从 channel 中取出并将东西放入 redis。
像这样:
c := make(chan Message, bufferLen)
go pollKafka(c)
go pushToRedis(c)
如果你想添加批处理,你可以添加一个从 kafka channel 轮询的中间阶段,并附加到一个 slice 直到 slice 已满,然后将该 slice 排入 redis 的 channel 。
如果这样的并发性不是目标,那么用 slice 替换代码中的 channel 可能会更容易。如果只有 1 个 goroutine 作用于一个对象,那么尝试使用 channel 并不是一个好主意。
关于与卡夫卡消费者一起去 channel ,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/50301067/
我想要的是能够在输入获得焦点或失去焦点时执行某些操作(两个事件)。 我尝试了以下方法,但这按事件单独工作(单独编码时):仅在焦点上,或仅在失去焦点时。 另外,我希望它尽可能跨平台(包括触摸设备),这是
我分别研究了TableView的Filtering和Pagination。 过滤: this帖子帮助我满足了我的需要 分页: this , this帖子也帮助了我 我想像这样将它们组合在一起: 详情-
我是 TDD 方法的新手,所以我想知道是否有人经历过这种机智可以启发我一点。我想获得一些关于如何一起使用 UML 和 TDD 方法的线索。 我已经习惯了:用 UML 设计 --> 生成骨架类(然后保持
我尝试使用入口点和 cmd 设置 Docker。 FROM debian:stretch RUN apt-get update && \ apt install gnupg ca-certificat
我想要一个 Class 对象,但我想强制它所代表的任何类扩展类 A 并实现接口(interface) B。 我能做到: Class 或者: Class 但我不能两者兼得。有办法做到这一点吗? 最佳答案
我是 Rubymine 的长期用户。 Rubymine 非常适合基于 html 的 Rails 应用程序,但我现在正在做更多的 SPA 客户端工作(例如 javascript/react)。我发现我真
我注意到我使用的某个脚本依赖于原型(prototype)。 (Lightbox 2) 它会与 jQuery 在同一页面上一起工作吗?有没有办法确保它们不冲突? 最佳答案 可以,但你需要采取 speci
我需要对表中显示的数据进行分页并通过 ajax 调用获取它 - 这是我通过使用具有以下配置的 dataTables 插件来完成的 - bServerSide : true; sAjaxSource :
我是 gtk 新手,所以想知道在 C 语言中归档和 gtk 是否可以一起使用?例如,我可以从 .txt 文件中读取,然后在相同的代码中使用 gtk 在标签或其他内容中显示它吗?如果是,怎么办? 谢谢!
有没有人设法得到Bck2Brwsr最近与 Java 8/JavaFX 8 一起工作?有没有兼容的机会?我找不到太多关于它的信息,也没有一个好的起点。使用给定的 Maven archetype我遇到了几
在我的应用程序中,用户通过 openid(与 stackoverflow 相同)登录/注销。 我想通过 oauth 向第三方应用程序开放我的应用程序。 如何创建我的 openid-consumer 应
我在启动和运行 Hibernate 和 Spring 时遇到一些问题。我有一个网络服务器项目,它使用了其他几个具有持久实体的项目。我遇到的问题是,对于存储在 WEB-INF/libs 内的另一个 ja
我有 @ControllerAdvice 类,它处理一组异常。我们还有一些其他异常,这些异常用 @ResponseStatus 注释进行注释。为了结合这两种方法,我们使用博客文章中描述的技术:http
我想在屏幕上使用进度条而不是 progressDialog。 我在我的 XML View 文件中插入了一个进度条,我想让它在加载时显示并在不加载时禁用它。 所以我使用的是可见的,但它发生了,所以其余的
CREATE TABLE `users` ( `id` int(11) AUTO_INCREMENT, `academicdegree` varchar(255),
IN() 中使用的查询返回:1, 2。然而,整个查询返回 0 行,这是不可能的,因为它们存在。我在这里做错了什么? SELECT DISTINCT li.auto_id FROM links
亲们, 我如何在使用 Jade 生成的表单上实现 jQuery 样式?我想做的是美化 表单并使它们可点击。我在 UI 方面很糟糕。期间。 我如何在表单上实现这个可选择的方法? http://jquer
按照目前的情况,这个问题不适合我们的问答形式。我们希望答案得到事实、引用或专业知识的支持,但这个问题可能会引发辩论、争论、投票或扩展讨论。如果您觉得这个问题可以改进并可能重新打开,visit the
我可以: auto o1 = new Content; 但不能: std::shared_ptr o1(new Content); std::unique_ptr o1(new Content); 我
关闭。这个问题需要更多focused .它目前不接受答案。 想改进这个问题吗? 更新问题,使其只关注一个问题 editing this post . 关闭 4 年前。 Improve this qu
我是一名优秀的程序员,十分优秀!