kafka介绍(卡夫卡的生平有哪些介绍)

本文目录
卡夫卡的生平有哪些介绍
卡夫卡(1883~1924年),奥地利小说家。出生在奥匈帝国统治下的布拉格,犹太血统。父亲是百货批发商。他18岁入布拉格大学学习文学和法律,毕业后主要从事保险业工作。1904年开始用德语写作,1917年起因患肺结核,辗转疗养。1923年迁居柏林,专事写作。但次年病情恶化,于同年7月3日病殁于维也纳。他的主要作品为4部短篇小说集和3部长篇小说。可惜生前大多未发表,3部长篇也均未写完。死后多年,其作品中的深刻思想逐渐被人认识,成为20世纪最伟大的现代主义作家之一。
卡夫卡的一生正值奥地利近代史上发生深刻社会变革的时期,他又深受尼采、柏格森哲学影响,对政治事件也一直抱旁观态度,故其作品大都用变形荒诞的形象和象征直觉的手法,表现被充满敌意的社会环境包围的孤立、绝望的个人。成为席卷欧洲的“现代人的困惑”的集中体现,并在欧洲掀起了一阵又一阵的“卡夫卡热”。其最著名的作品有借小动物防备敌害的胆战心理,擅长表现资本主义社会小人物时刻难以自保的精神状态和在充满敌意的环境中的孤立绝望的情绪。
卡夫卡简介
卡夫卡简介弗兰兹· 卡夫卡( (Franz Kafka 1883 ~1924)奥地利小说家。出生犹太商人家庭,)奥地利小说家。出生犹太商人家庭,18岁入布拉格大学学习文学和法律,岁入布拉格大学学习文学和法律,1904年开始写作,主要作品为年开始写作,主要作品为4 部短篇小说集和3部长篇小说。可惜生前大多未发表,3部长篇也均未写完。卡夫卡是欧洲著名的表现主义作家。他生活在奥匈帝国行将崩溃的时代,又深受尼采、柏格森哲学影响,对政治事件也一直抱旁观态度,故其作品大都用变形荒诞的形象和象征直觉的手法,表现被充满敌意的社会环境所包围的孤立、绝望的个人。部长篇也均未写完。卡夫卡是欧洲著名的表现主义作家。他生活在奥匈帝国行将崩溃的时代,又深受尼采、柏格森哲学影响,对政治事件也一直抱旁观态度卡夫卡简介,故其作品大都用变形荒诞的形象和象征直觉的手法,表现被充满敌意的社会环境所包围的孤立、绝望的个人。卡夫卡生平介绍• 弗兰兹·卡夫卡(Franz Kafka,1883年7月3日—1924年6月3日),20世纪最有影响力的德语小说家。• 弗兰兹· 卡夫卡• 文笔明净而想像奇诡,常采用寓言体,背后的寓意言人人殊,暂无(或永无)定论。
别开生面的手法,令二十世纪各个写作流派纷纷追认其为先驱。• 卡夫卡是奥地利人,他是西方现代派文学的宗师和探险者,他的创作风格是表现主义,是表现主义作家中创作上最有成就者。他生活和创作的主要时期是在一战前后,当时,经济萧条,社会腐败,人民穷困,这一切使得卡夫卡终生生活在痛苦与孤独之中。于是,对社会的陌生感,孤独感与恐惧感,成了他创作的永恒主题。美国诗人奥登评价卡夫卡时说:“卡夫卡对我们至关重要,因为他的困境就是现代人的困境。• 卡夫卡一生的作品并不多,但对后世文学的影响却是极为深远的。美国诗人奥登认为:“他与我们时代的关系最近似但丁、莎士比亚、歌德与他们时代的关系。”卡夫卡的小说揭示了一种荒诞的充满非理性色彩的景象,个人式的、忧郁的、孤独的情绪,运用的是象征式的手法。三四十年代的超现实主义余党视之为同仁,四五十年代的荒诞派以之为先驱,六十年代的美卡夫卡一生的作品并不多,但对后世文学的影响却是极为深远的。美国诗人奥登认为:“他与我们时代的关系最近似但丁、莎士比亚、歌德与他们时代的关系。”卡夫卡的小说揭示了一种荒诞的充满非理性色彩的景象,个人式的、忧郁的、孤独的情绪,运用的是象征式的手法。
三四十年代的超现实主义余党视之为同仁,四五十年代的荒诞派以之为先驱,六十年代的美• 弗兰兹· 卡夫卡• 国”黑色幽默“奉之为典范。艺术特点卡夫卡被认为是现代派文学的鼻祖,是表现主义文学的先驱,其作品主题曲折晦涩,情节支离破碎,思路不连贯,跳跃性很大,语言的象征意义很强,这给阅读和理解他的作品带来了一定的困难。卡夫卡的作品难读,连母语是德语的读者也觉得读懂这些作品不是件容易的事,但他那独到的认识,深刻的批判,入木三分的描写,都深深地吸引着我们,只要你能读进去,只要你能摸到作品的脉络,定会获益匪浅。下点功夫读一读卡夫卡是值得的。卡夫卡笔下描写的都是生活在下层的小人物,他们在这充满矛盾、扭曲变形的世界里惶恐,不安,孤独,迷惘,遭受压迫而不敢反抗,也无力反抗,向往明天又看不到出路。看到他为我们描绘出的一幅幅画卷我们会感到一阵阵震惊和恐惧,因为他仿佛在为人类的明天敲起阵阵急促的警钟,他为人类的未来担忧。每位读者在读卡夫卡时都会有自己的感触、理解、认识、联想,但我们希望读者不要迷惘在他所描绘的迷惘中。卡夫卡与中国阿根廷杰出的小说家,诗人博尔赫斯(Jorge Luis Borges)是首个将卡夫卡小说译为西班牙文的人,他在一篇文章《)是首个将卡夫卡小说译为西班牙文的人,他在一篇文章《Kafka y sus precursores》中替卡夫卡追宗认祖,其中一人是韩退之,全因他写过《获麟解》这篇寓言。
卡夫卡读过一些中国文学的德译本,他在》中替卡夫卡追宗认祖,其中一人是韩退之,全因他写过《获麟解》这篇寓言。卡夫卡读过一些中国文学的德译本,他在1912年写信给当时的未婚妻,引用了袁枚一首不太高明的诗:年写信给当时的未婚妻,引用了袁枚一首不太高明的诗:《寒夜》寒夜读书忘却眠锦衾香尽炉无烟美人含怒夺灯去问郎知是几更天《寒夜》寒夜读书忘却眠锦衾香尽炉无烟美人含怒夺灯去问郎知是几更天卡夫卡读过中国古代哲学家的文学著作,有《南华经》《论语》《道德经》等。卡夫卡偏爱研究道家,他说:“在孔子的《论语》里,起初人们还站在坚实的大地上,但到后来书里的内容越来越虚无缥缈,让读者不可捉摸。老子的格言是坚硬的核桃,我被它们陶醉了,但是它们的核心对我依然紧锁着。我反复读了好多遍,然后我却发现,就像小孩玩彩色玻璃球游戏那样,我让这些格言从一个思想角落滑到另一个思想角落,而丝毫没有前进。通过这些格言玻璃球,我其实只发现了我的思想槽非常浅,无法包容老子的玻璃球。这是令人沮丧的发现,于是我停止了玻璃球戏。”卡夫卡读过中国古代哲学家的文学著作,有《南华经》《论语》《道德经》等。卡夫卡偏爱研究道家,他说:“在孔子的《论语》里,起初人们还站在坚实的大地上,但到后来书里的内容越来越虚无缥缈,让读者不可捉摸。
老子的格言是坚硬的核桃,我被它们陶醉了,但是它们的核心对我依然紧锁着。我反复读了好多遍,然后我却发现,就像小孩玩彩色玻璃球游戏那样,我让这些格言从一个思想角落滑到另一个思想角落,而丝毫没有前进。通过这些格言玻璃球,我其实只发现了我的思想槽非常浅,无法包容老子的玻璃球。这是令人沮丧的发现,于是我停止了玻璃球戏。”作品列举生前出版的单行本《判决》(Das Urteil) 《火夫》(或译《司炉》)(Der Heizer) 《变形记》(Die Verwandlung) 《在流放地》(In der Strafkolonie)生前出版的集子《观察》生前出版的集子《观察》(Betrachtung) 《乡村医生》(Ein Landarzt) 《饥饿艺术家》(Ein Hungerkünstler)生前出版的小说(未结集)《与祈祷者的对话》(Gespräch mit dem Beter) 《与醉汉的对话》(Gespräch mit dem Betrunkenen) 《巨响》《桶骑士》《巨响》《桶骑士》(Der Kübelreiter)遗作(长篇小说)《失踪者》(Der Verschollene) 【一名《美国》【一名《美国》(Amerika)】《审判》(或译《诉讼》)】《审判》(或译《诉讼》)(Der Prozeß) 《城堡》(Das Schloß)作品介绍《审判》卡夫卡的作品据说是西方现代文学作品中最难读的一种,最主要的原因是他的小说负担的内容太多。
根据我读卡夫卡的经验,在阅读卡夫卡之前必须有两方面的前理解准备:一是对西方文明的源----所谓“二希”:古希腊的哲学和希伯莱的宗教,和流---所谓康德之后乃至尼采之后--的嬗变有一个了解。二是对卡夫卡个人性情的了解。据说巴尔扎克在他的手杖上刻着一句话:我粉碎一切障碍,卡夫卡反其意而用之,说一切障碍都在粉碎我。卡夫卡的理想是做一个地窖隐士,在昏暗的地窖之中不受打扰地用写作滋润自己的灵魂。还应该注意的是卡夫卡与他父亲的关系。卡夫卡的长篇小说《审判》是卡夫卡形成自己风格的第一部长篇小说。《审判》的创作与卡夫卡订婚-解除婚约--又订婚的经历重合。小说讲的是银行高级职员约瑟夫.K在三十岁生日那天突然被一群神秘的黑衣人宣布有罪,但是他又是自由的,于是他开始了艰难的上诉之路,但是毫无结果,在三十一岁生日那天被秘密处决。首先我们关注小说的开头。K在三十岁生日那天早晨醒来时突然被宣布有罪。生日意味着什么?生日意味着我们出生了,但是出生并不是我们的意志,我们是被动的,我们被出生(be born),我们没有经过自己的同意被抛到这个世界上。如果我们考虑到卡夫卡对世界的悲观态度,那么我们可以说出生是一种抛弃。
但丁在《神曲》的开头说,他在人生的中途,30岁时步入歧路,前有狼,后有狮,在诗人维吉尔的引导下游历了地狱、炼狱,在女友的引导下游历了天堂。因此30岁是个很有意思的分界,中国的孔圣也说三十而立。三十岁似乎是一个人智性觉醒的时期。如果说三十岁之前的人是做为一个自在的人而存在的话,三十岁以后的• 人做为一个自为的人而存在。在这个“新生”的早晨,K被宣布“有罪”。在被宣布有罪之后,由于早餐被黑衣人享用了,K只好找点东西当早餐,他先是找到了一只苹果,然后又喝了点酒。请注意在文本中卡夫卡对苹果的形• 容:“漂亮的”,这是在小说开头灰暗的文本中间唯一一个温暖的词。苹果而不是其他的水果让人想起《圣经》中的相关描述,苹果是知识之树上的果子,人类之祖因受到蛇的诱惑吃了这个果子后被宣布有罪而赶出了伊甸园。因此,苹果代表理性的觉醒,是人对自己无辜的一种自觉。吃完苹果后K又喝了点酒,这不禁让人想起尼采,K不但是康德以后--信仰的上帝被杀死以后,而且是尼采之后--道德的上帝被杀死以后的人,是自知自己的无辜而要求上诉的人。
本来在上帝的法庭上没有上诉的可能,末日审判是绝对的终审判决,古人的罪是自觉的罪,是道德堕落意义上对上帝所犯的罪。现代人的罪感发生了变化,现代人所感到的是生成的无辜,是某种自然意义上的欠缺,是面对生命的偶然时的终究意难平。正是这种关于罪的感觉的颠转,造成K上诉的前提。整部小说因此很象《苏格拉底的申辩》,是在上帝面前对生存感觉发生变化的人类所做的辩护,或者说在上帝的法庭上辩白人生成的无辜。但是小说整个阴沉的格调显示了这种在神义论面前为人义论辩护的艰难。• 值得注意的是,整部小说的倒数第二章是K与教士的对话,然后,在最后一章,K在31岁生日那天被秘密处决。与生一样,远离上帝,现代人的死也变成了一种”横死”。古人一般都相信人死后会变成鬼,鬼者,归也。死亡是一种回归,对有永生信仰的人来说,死亡是一种判决,或者入地狱,或者进天堂。但是对于祛魅后的现代人来说,死亡没有意义,死亡是诸种偶然性中的一种,死亡不再是一种判决,死亡下面是无尽的虚无,死亡是对人生无意义的最深佐证. 小说的最后,K仍然想着是否有改判的可能,秘密处死是不是必然的命运,黑夜里对面楼里的灯光昏暗,黑衣人的刀 *** K的胸膛,并转动了两下,灯光逐渐模糊. 《城堡》• 《城堡》是卡夫卡最后一部长篇小说,没有完成。
小说写的是主人公K为了进入城堡而努力的故事。一个冬夜,K经过长途跋涉来到了城堡所属的一个村庄,投宿在一个乡村客店里。按照规定没有城堡的许可谁都不能在村子里过夜。K自称是土地测量员给城堡工作,由于没有任何证据又遭到严厉的盘查。客栈用电话向城堡查询这件事情以后得到肯定的答复,K才被允许留下来过夜。其实,城堡根本没有聘请K来工作,却承认了他并给他派了两个助手,只是始终不允许他进入城堡。尽管城堡就在近在咫尺的小山上,却是永远可遇不可求的,他永远也走不到那里。为了能进去,他有意勾引了城堡办公厅主任的情人佛利达,之后发生的一切,佛利达的猜忌,给K送信的巴纳巴斯家的不幸,与克拉姆秘书在贵宾室几经波折的会面,佛利达处于猜忌和妒忌的私奔同居,一系列这些事情我都觉得很莫名,简直没有前因后果,小说写到这个时候就停止了。《城堡》凝聚了他长久的人生思考,表达了他对社会,对亲情,对爱情卡夫卡简介,对生计等所有的一切的理解,虽然似懂非懂的读完了这本书,但它还是触动了已积淀很久,快被遗忘的心灵感受。• 我觉得,读卡夫卡就是因为他的那种隐晦的比喻,一旦领会了他的象征,感触颇多。
kafka原理分析
作为一款典型的消息中间件产品,kafka系统仍然由producer、broker、consumer三部分组成。kafka涉及的几个常用概念和组件简单介绍如下:
当consumer group的状态发生变化(如有consumer故障、增减consumer成员等)或consumer group消费的topic状态发生变化(如增加了partition,消费的topic发生变化),kafka集群会自动调整和重新分配consumer消费的partition,这个过程就叫做rebalance(再平衡)。
__consumer_offsets是kafka集群自己维护的一个特殊的topic,它里面存储的是每个consumer group已经消费了每个topic partition的offset。__consumer_offsets中offset消息的key由group id,topic name,partition id组成,格式为 {topic name}-${partition id},value值就是consumer提交的已消费的topic partition offset值。__consumer_offsets的分区数和副本数分别由offsets.topic.num.partitions(默认值为50)和offsets.topic.replication.factor(默认值为1)参数配置。我们通过公式 hash(group id) % offsets.topic.num.partitions 就可以计算出指定consumer group的已提交offset存储的partition。由于consumer group提交的offset消息只有最后一条消息有意义,所以__consumer_offsets是一个compact topic,kafka集群会周期性的对__consumer_offsets执行compact操作,只保留最新的一次提交offset。
group coordinator运行在kafka某个broker上,负责consumer group内所有的consumer成员管理、所有的消费的topic的partition的消费关系分配、offset管理、触发rebalance等功能。group coordinator管理partition分配时,会指定consumer group内某个consumer作为group leader执行具体的partition分配任务。存储某个consumer group已提交offset的__consumer_offsets partition leader副本所在的broker就是该consumer group的协调器运行的broker。
跟大多数分布式系统一样,集群有一个master角色管理整个集群,协调集群中各个成员的行为。kafka集群中的controller就相当于其它分布式系统的master,用来负责集群topic的分区分配,分区leader选举以及维护集群的所有partition的ISR等集群协调功能。集群中哪个borker是controller也是通过一致性协议选举产生的,2.8版本之前通过zookeeper进行选主,2.8版本后通过kafka raft协议进行选举。如果controller崩溃,集群会重新选举一个broker作为新的controller,并增加controller epoch值(相当于zookeeper ZAB协议的epoch,raft协议的term值)
当kafka集群新建了topic或为一个topic新增了partition,controller需要为这些新增加的partition分配到具体的broker上,并把分配结果记录下来,供producer和consumer查询获取。
因为只有partition的leader副本才会处理producer和consumer的读写请求,而partition的其他follower副本需要从相应的leader副本同步消息,为了尽量保证集群中所有broker的负载是均衡的,controller在进行集群全局partition副本分配时需要使partition的分布情况是如下这样的:
在默认情况下,kafka采用轮询(round-robin)的方式分配partition副本。由于partition leader副本承担的流量比follower副本大,kafka会先分配所有topic的partition leader副本,使所有partition leader副本全局尽量平衡,然后再分配各个partition的follower副本。partition第一个follower副本的位置是相应leader副本的下一个可用broker,后面的副本位置依此类推。
举例来说,假设我们有两个topic,每个topic有两个partition,每个partition有两个副本,这些副本分别标记为1-1-1,1-1-2,1-2-1,1-2-2,2-1-1,2-1-2,2-2-1,2-2-2(编码格式为topic-partition-replia,编号均从1开始,第一个replica是leader replica,其他的是follower replica)。共有四个broker,编号是1-4。我们先对broker按broker id进行排序,然后分配leader副本,最后分配foller副本。
1)没有配置broker.rack的情况
现将副本1-1-1分配到broker 1,然后1-2-1分配到broker 2,依此类推,2-2-1会分配到broker 4。partition 1-1的leader副本分配在broker 1上,那么下一个可用节点是broker 2,所以将副本1-1-2分配到broker 2上。同理,partition 1-2的leader副本分配在broker 2上,那么下一个可用节点是broker 3,所以将副本1-1-2分配到broker 3上。依此类推分配其他的副本分片。最后分配的结果如下图所示:
2)配置了broker.rack的情况
假设配置了两个rack,broker 1和broker 2属于Rack 1,broker 3和broker 4属于Rack 2。我们对rack和rack内的broker分别排序。然后先将副本1-1-1分配到Rack 1的broker 1,然后将副本1-2-1分配到下一个Rack的第一个broker,即Rack 2的broker 3。其他的parttition leader副本依此类推。然后分配follower副本,partition 1-1的leader副本1-1-1分配在Rack 1的broker上,下一个可用的broker是Rack 2的broker 3,所以分配到broker 3上,其他依此类推。最后分配的结果如下图所示:
kafka除了按照集群情况自动分配副本,也提供了reassign工具人工分配和迁移副本到指定broker,这样用户可以根据集群实际的状态和各partition的流量情况分配副本
kafka集群controller的一项功能是在partition的副本中选择一个副本作为leader副本。在topic的partition创建时,controller首先分配的副本就是leader副本,这个副本又叫做preference leader副本。
当leader副本所在broker失效时(宕机或网络分区等),controller需要为在该broker上的有leader副本的所有partition重新选择一个leader,选择方法就是在该partition的ISR中选择第一个副本作为新的leader副本。但是,如果ISR成员只有一个,就是失效的leader自身,其余的副本都落后于leader怎么办?kafka提供了一个unclean.leader.election配置参数,它的默认值为true。当unclean.leader.election值为true时,controller还是会在非ISR副本中选择一个作为leader,但是这时候使用者需要承担数据丢失和数据不一致的风险。当unclean.leader.election值为false时,则不会选择新的leader,该partition处于不可用状态,只能恢复失效的leader使partition重新变为可用。
当preference leader失效后,controller重新选择一个新的leader,但是preference leader又恢复了,而且同步上了新的leader,是ISR的成员,这时候preference leader仍然会成为实际的leader,原先的新leader变为follower。因为在partition leader初始分配时,使按照集群副本均衡规则进行分配的,这样做可以让集群尽量保持平衡。
为了保证topic的高可用,topic的partition往往有多个副本,所有的follower副本像普通的consumer一样不断地从相应的leader副本pull消息。每个partition的leader副本会维护一个ISR列表存储到集群信息库里,follower副本成为ISR成员或者说与leader是同步的,需要满足以下条件:
1)follower副本处于活跃状态,与zookeeper(2.8之前版本)或kafka raft master之间的心跳正常
2)follower副本最近replica.lag.time.max.ms(默认是10秒)时间内从leader同步过最新消息。需要注意的是,一定要拉取到最新消息,如果最近replica.lag.time.max.ms时间内拉取过消息,但不是最新的,比如落后follower在追赶leader过程中,也不会成为ISR。
follower在同步leader过程中,follower和leader都会维护几个参数,来表示他们之间的同步情况。leader和follower都会为自己的消息队列维护LEO(Last End Offset)和HW(High Watermark)。leader还会为每一个follower维护一个LEO。LEO表示leader或follower队列写入的最后一条消息的offset。HW表示的offset对应的消息写入了所有的ISR。当leader发现所有follower的LEO的最小值大于HW时,则会增加HW值到这个最小值LEO。follower拉取leader的消息时,同时能获取到leader维护的HW值,如果follower发现自己维护的HW值小于leader发送过来的HW值,也会增加本地的HW值到leader的HW值。这样我们可以得到一个不等式: follower HW 《= leader HW 《= follower LEO 《= leader LEO 。HW对应的log又叫做committed log,consumer消费partititon的消息时,只能消费到offset值小于或等于HW值的消息的,由于这个原因,kafka系统又称为分布式committed log消息系统。
kafka的消息内容存储在log.dirs参数配置的目录下。kafka每个partition的数据存放在本地磁盘log.dirs目录下的一个单独的目录下,目录命名规范为 ${topicName}-${partitionId} ,每个partition由多个LogSegment组成,每个LogSegment由一个数据文件(命名规范为: {baseOffset}.index)和一个时间戳索引文件(命名规范为:${baseOffset}.timeindex)组成,文件名的baseOffset就是相应LogSegment中第一条消息的offset。.index文件存储的是消息的offset到该消息在相应.log文件中的偏移,便于快速在.log文件中快速找到指定offset的消息。.index是一个稀疏索引,每隔一定间隔大小的offset才会建立相应的索引(比如每间隔10条消息建立一个索引)。.timeindex也是一个稀疏索引文件,这样可以根据消息的时间找到对应的消息。
可以考虑将消息日志存放到多个磁盘中,这样多个磁盘可以并发访问,增加消息读写的吞吐量。这种情况下,log.dirs配置的是一个目录列表,kafka会根据每个目录下partition的数量,将新分配的partition放到partition数最少的目录下。如果我们新增了一个磁盘,你会发现新分配的partition都出现在新增的磁盘上。
kafka提供了两个参数log.segment.bytes和log.segment.ms来控制LogSegment文件的大小。log.segment.bytes默认值是1GB,当LogSegment大小达到log.segment.bytes规定的阈值时,kafka会关闭当前LogSegment,生成一个新的LogSegment供消息写入,当前供消息写入的LogSegment称为活跃(Active)LogSegment。log.segment.ms表示最大多长时间会生成一个新的LogSegment,log.segment.ms没有默认值。当这两个参数都配置了值,kafka看哪个阈值先达到,触发生成新的LogSegment。
kafka还提供了log.retention.ms和log.retention.bytes两个参数来控制消息的保留时间。当消息的时间超过了log.retention.ms配置的阈值(默认是168小时,也就是一周),则会被认为是过期的,会被kafka自动删除。或者是partition的总的消息大小超过了log.retention.bytes配置的阈值时,最老的消息也会被kafka自动删除,使相应partition保留的总消息大小维持在log.retention.bytes阈值以下。这个地方需要注意的是,kafka并不是以消息为粒度进行删除的,而是以LogSegment为粒度删除的。也就是说,只有当一个LogSegment的最后一条消息的时间超过log.retention.ms阈值时,该LogSegment才会被删除。这两个参数都配置了值时,也是只要有一个先达到阈值,就会执行相应的删除策略
当我们使用KafkaProducer向kafka发送消息时非常简单,只要构造一个包含消息key、value、接收topic信息的ProducerRecord对象就可以通过KafkaProducer的send()向kafka发送消息了,而且是线程安全的。KafkaProducer支持通过三种消息发送方式
KafkaProducer客户端虽然使用简单,但是一条消息从客户端到topic partition的日志文件,中间需要经历许多的处理过程。KafkaProducer的内部结构如下所示:
从图中可以看出,消息的发送涉及两类线程,一类是调用KafkaProducer.send()方法的应用程序线程,因为KafkaProducer.send()是多线程安全的,所以这样的线程可以有多个;另一类是与kafka集群通信,实际将消息发送给kafka集群的Sender线程,当我们创建一个KafkaProducer实例时,会创建一个Sender线程,通过该KafkaProducer实例发送的所有消息最终通过该Sender线程发送出去。RecordAccumulator则是一个消息队列,是应用程序线程与Sender线程之间消息传递的桥梁。当我们调用KafkaProducer.send()方法时,消息并没有直接发送出去,只是写入了RecordAccumulator中相应的队列中,最终需要Sender线程在适当的时机将消息从RecordAccumulator队列取出来发送给kafka集群。
消息的发送过程如下:
在使用KafkaConsumer实例消费kafka消息时,有一个特性我们要特别注意,就是KafkaConsumer不是多线程安全的,KafkaConsumer方法都在调用KafkaConsumer的应用程序线程中运行(除了consumer向kafka集群发送的心跳,心跳在一个专门的单独线程中发送),所以我们调用KafkaConsumer的所有方法均需要保证在同一个线程中调用,除了KafkaConsumer.wakeup()方法,它设计用来通过其它线程向consumer线程发送信号,从而终止consumer执行。
跟producer一样,consumer要与kafka集群通信,消费kafka消息,首先需要获取消费的topic partition leader replica所在的broker地址等信息,这些信息可以通过向kafka集群任意broker发送Metadata请求消息获取。
我们知道,一个consumer group有多个consumer,一个topic有多个partition,而且topic的partition在同一时刻只能被consumer group内的一个consumer消费,那么consumer在消费partition消息前需要先确定消费topic的哪个partition。partition的分配通过group coordinator来实现。基本过程如下:
我们可以通过实现接口org.apache.kafka.clients.consumer.internals.PartitionAssignor自定义partition分配策略,但是kafka已经提供了三种分配策略可以直接使用。
partition分配完后,每个consumer知道了自己消费的topic partition,通过metadata请求可以获取相应partition的leader副本所在的broker信息,然后就可以向broker poll消息了。但是consumer从哪个offset开始poll消息?所以consumer在第一次向broker发送FetchRequest poll消息之前需要向Group Coordinator发送OffsetFetchRequest获取消费消息的起始位置。Group Coordinator会通过key {topic}-${partition}查询 __consumer_offsets topic中是否有offset的有效记录,如果存在,则将consumer所属consumer group最近已提交的offset返回给consumer。如果没有(可能是该partition是第一次分配给该consumer group消费,也可能是该partition长时间没有被该consumer group消费),则根据consumer配置参数auto.offset.reset值确定consumer消费的其实offset。如果auto.offset.reset值为latest,表示从partition的末尾开始消费,如果值为earliest,则从partition的起始位置开始消费。当然,consumer也可以随时通过KafkaConsumer.seek()方法人工设置消费的起始offset。
kafka broker在收到FetchRequest请求后,会使用请求中topic partition的offset查一个skiplist表(该表的节点key值是该partition每个LogSegment中第一条消息的offset值)确定消息所属的LogSegment,然后继续查LogSegment的稀疏索引表(存储在.index文件中),确定offset对应的消息在LogSegment文件中的位置。为了提升消息消费的效率,consumer通过参数fetch.min.bytes和max.partition.fetch.bytes告诉broker每次拉取的消息总的最小值和每个partition的最大值(consumer一次会拉取多个partition的消息)。当kafka中消息较少时,为了让broker及时将消息返回给consumer,consumer通过参数fetch.max.wait.ms告诉broker即使消息大小没有达到fetch.min.bytes值,在收到请求后最多等待fetch.max.wait.ms时间后,也将当前消息返回给consumer。fetch.min.bytes默认值为1MB,待fetch.max.wait.ms默认值为500ms。
为了提升消息的传输效率,kafka采用零拷贝技术让内核通过DMA把磁盘中的消息读出来直接发送到网络上。因为kafka写入消息时将消息写入内存中就返回了,如果consumer跟上了producer的写入速度,拉取消息时不需要读磁盘,直接从内存获取消息发送出去就可以了。
为了避免发生再平衡后,consumer重复拉取消息,consumer需要将已经消费完的消息的offset提交给group coordinator。这样发生再平衡后,consumer可以从上次已提交offset出继续拉取消息。
kafka提供了多种offset提交方式
partition offset提交和管理对kafka消息系统效率来说非常关键,它直接影响了再平衡后consumer是否会重复拉取消息以及重复拉取消息的数量。如果offset提交的比较频繁,会增加consumer和kafka broker的消息处理负载,降低消息处理效率;如果offset提交的间隔比较大,再平衡后重复拉取的消息就会比较多。还有比较重要的一点是,kafka只是简单的记录每次提交的offset值,把最后一次提交的offset值作为最新的已提交offset值,作为再平衡后消息的起始offset,而什么时候提交offset,每次提交的offset值具体是多少,kafka几乎不关心(这个offset对应的消息应该存储在kafka中,否则是无效的offset),所以应用程序可以先提交3000,然后提交2000,再平衡后从2000处开始消费,决定权完全在consumer这边。
kafka中的topic partition与consumer group中的consumer的消费关系其实是一种配对关系,当配对双方发生了变化时,kafka会进行再平衡,也就是重新确定这种配对关系,以提升系统效率、高可用性和伸缩性。当然,再平衡也会带来一些负面效果,比如在再平衡期间,consumer不能消费kafka消息,相当于这段时间内系统是不可用的。再平衡后,往往会出现消息的重复拉取和消费的现象。
触发再平衡的条件包括:
需要注意的是,kafka集群broker的增减或者topic partition leader重新选主这类集群状态的变化并不会触发在平衡
有两种情况与日常应用开发比较关系比较密切:
consumer在调用subscribe()方法时,支持传入一个ConsumerRebalanceListener监听器,ConsumerRebalanceListener提供了两个方法,onPartitionRevoked()方法在consumer停止消费之后,再平衡开始之前被执行。可以发现,这个地方是提交offset的好时机。onPartitonAssigned()方法则会在重新进行partition分配好了之后,但是新的consumer还未消费之前被执行。
我们在提到kafka时,首先想到的是它的吞吐量非常大,这也是很多人选择kafka作为消息传输组件的重要原因。
以下是保证kafka吞吐量大的一些设计考虑:
但是kafka是不是总是这么快?我们同时需要看到kafka为了追求快舍弃了一些特性:
所以,kafka在消息独立、允许少量消息丢失或重复、不关心消息顺序的场景下可以保证非常高的吞吐量,但是在需要考虑消息事务、严格保证消息顺序等场景下producer和consumer端需要进行复杂的考虑和处理,可能会比较大的降低kafka的吞吐量,例如对可靠性和保序要求比较高的控制类消息需要非常谨慎的权衡是否适合使用kafka。
我们通过producer向kafka集群发送消息,总是期望消息能被consumer成功消费到。最不能忍的是producer收到了kafka集群消息写入的正常响应,但是consumer仍然没有消费到消息。
kafka提供了一些机制来保证消息的可靠传递,但是有一些因素需要仔细权衡考虑,这些因素往往会影响kafka的吞吐量,需要在可靠性与吞吐量之间求得平衡:
kafka只保证partition消息顺序,不保证topic级别的顺序,而且保证的是partition写入顺序与读取顺序一致,不是业务端到端的保序。
如果对保序要求比较高,topic需要只设置一个partition。这时可以把参数max.in.flight.requests.per.connection设置为1,而retries设置为大于1的数。这样即使发生了可恢复型错误,仍然能保证消息顺序,但是如果发生不可恢复错误,应用层进行重试的话,就无法保序了。也可以采用同步发送的方式,但是这样也极大的降低了吞吐量。如果消息携带了表示顺序的字段,可以在接收端对消息进行重新排序以保证最终的有序。

更多文章:
在from子句中可以出现(如何在from 子句中嵌套查询下面的语句在access中出错!)
2026年10月11日 05:20
countif函数统计个数怎么用(countif函数怎么用 详解Excel中countif函数的使用方法)
2026年10月11日 03:30
正则匹配数字之前的字符(正则表达式如何匹配前面是数字、中间是“/”、后面也是数字,就像2/3专业的模式)
2026年10月11日 03:00
orlnsertbootmediinselected(我电脑开机显示这个是什么意思or insert boot media in select)
2026年10月10日 23:00
display的用法(display是什么意思 详解display的含义和用法)
2026年10月10日 22:00
html全部居中代码(怎么让网页居中显示,html如何让网页居中)
2026年10月10日 21:10





