【问题标题】:Doesn't receive a message from a channel没有收到来自频道的消息
【发布时间】:2018-09-02 14:16:54
【问题描述】:

编辑:

在我添加了我正在使用的文件的一小部分(7 GB)并尝试运行程序后,我可以看到:

fatal error: all goroutines are asleep - deadlock!

goroutine 1 [chan receive]:
main.main()
    /media/developer/golang/manual/examples/sp/v2/sp.v2.go:71 +0x4a9
exit status 2

情况:

我是 GO 的新手,如果我的问题真的很简单,我很抱歉。

我正在尝试流式传输 xml 文件,拆分文档,然后在不同的 GO 例程中解析它们。

我正在使用的 XML 文件示例:

<?xml version="1.0" encoding="UTF-8"?>
<osm version="0.6" generator="CGImap 0.0.2">
    <relation id="56688" user="kmvar" uid="56190" visible="true" version="28" changeset="6947637" timestamp="2011-01-12T14:23:49Z">
        <member type="node" ref="294942404" role=""/>
        <member type="node" ref="364933006" role=""/>
        <tag k="name" v="Küstenbus Linie 123"/>
        <tag k="network" v="VVW"/>
        <tag k="route" v="bus"/>
        <tag k="type" v="route"/>
    </relation>
    <relation id="98367" user="jdifh" uid="92834" visible="true" version="28" changeset="6947637" timestamp="2011-01-12T14:23:49Z">
        <member type="node" ref="294942404" role=""/>
        <member type="way" ref="4579143" role=""/>
        <member type="node" ref="249673494" role=""/>
        <tag k="name" v="Küstenbus Linie 123"/>
        <tag k="network" v="VVW"/>
        <tag k="operator" v="Regionalverkehr Küste"/>
        <tag k="ref" v="123"/>
    </relation>
    <relation id="72947" user="schsu" uid="92374" visible="true" version="28" changeset="6947637" timestamp="2011-01-12T14:23:49Z">
        <member type="node" ref="294942404" role=""/>
        <tag k="name" v="Küstenbus Linie 123"/>
        <tag k="type" v="route"/>
    </relation>
    <relation id="93742" user="doiff" uid="61731" visible="true" version="28" changeset="6947637" timestamp="2011-01-12T14:23:49Z">
        <member type="node" ref="294942404" role=""/>
        <member type="node" ref="364933006" role=""/>
        <tag k="route" v="bus"/>
        <tag k="type" v="route"/>
    </relation>
</osm>

我有这段代码:

package main

import (
  "encoding/xml"
  "bufio"
  "fmt"
  "os"
  "io"
)

type RS struct {
  I string `xml:"id,attr"`
  M []struct {
    I string `xml:"ref,attr"`
    T string `xml:"type,attr"`
    R string `xml:"role,attr"`
  } `xml:"member"`
  T []struct {
    K string `xml:"k,attr"`
    V string `xml:"v,attr"`
  } `xml:"tag"`
}

func main() {
  p1D, err := os.Open("/media/developer/Transcend/osm/xml/relations.xml")

  if err != nil {
    fmt.Println(err)
    os.Exit(1)
  }

  defer p1D.Close()

  reader := bufio.NewReader(p1D)

  var count int32
  var element string

  channel := make(chan RS) // channel

  for {
    p2Rc, err := reader.ReadSlice('\n')
    if err != nil {
      if err == io.EOF {
        break
      } else {
        fmt.Println(err)
        os.Exit(1)
      }
    }

    var p2Rs = string(p2Rc)

    if p2Rc[2] == 114 {
      count++

      if (count != 1) {
        go parseRelation(element, channel)
      }

      element = ""
      element += p2Rs
    } else {
      element += p2Rs
    }
  }

  for i := 0; i < 5973974; i++ {
    fmt.Println(<- channel)
  }
}

func parseRelation(p1E string, channel chan RS) {
  var rs RS
  xml.Unmarshal([]byte(p1E), &rs)

  channel <- rs
}

它应该打印每个结构,但我什么也没看到。程序只是挂起。

我测试了流媒体和拆分器(刚刚在函数 parseRelation 中添加了fmt.Println(rs),然后将消息发送到通道中)。我可以看到结构。所以,问题在于发送和接收消息。

问题:

我不知道如何解决这个问题。尝试更改频道中消息的类型(从RS 到string)并且每次只发送一个字符串。但这也无济于事(我什么也看不见)

【问题讨论】:

  • 5973974这个数字从何而来?是某人的电话号码吗?
  • @DietrichEpp 就是文件中xml文档的数量。
  • 请发布您正在使用的 xml 示例,它将帮助我们检查数据是否未编组。
  • @Himanshu 刚刚添加了xml文件的例子。现在,当我使用 xml 文件的小示例运行程序时,我可以看到错误 fatal error: all goroutines are asleep - deadlock!

标签: go channel goroutine


【解决方案1】:

首先,让我们解决这个问题:您不能逐行解析 XML。您很幸运,您的文件恰好是每行一个标签,但这不能被视为理所当然。您必须解析整个 XML 文档。

通过逐行处理,您试图将&lt;tag&gt; 和&lt;member&gt; 推入为&lt;relation&gt; 设计的结构中。相反,请使用 xml.NewDecoder 并让它为您处理文件。

package main

import (
    "encoding/xml"
    "fmt"
    "os"
    "log"
)

type Osm struct {
    XMLName     xml.Name    `xml:"osm"`
    Relations   []Relation  `xml:"relation"`
}
type Relation struct {
    XMLName     xml.Name    `xml:"relation"`
    ID          string      `xml:"id,attr"`
    User        string      `xml:"user,attr"`
    Uid         string      `xml:"uid,attr"`
    Members     []Member    `xml:"member"`
    Tags        []Tag       `xml:"tag"`
}
type Member struct {
    XMLName     xml.Name    `xml:"member"`
    Ref         string      `xml:"ref,attr"`
    Type        string      `xml:"type,attr"`
    Role        string      `xml:"role,attr"`
}
type Tag struct {
    XMLName     xml.Name    `xml:"tag"`
    Key         string      `xml:"k,attr"`
    Value       string      `xml:"v,attr"`
}

func main() {
    reader, err := os.Open("test.xml")
    if err != nil {
        log.Fatal(err)
    }
    defer reader.Close()

    decoder := xml.NewDecoder(reader)

    osm := &Osm{}
    err = decoder.Decode(&osm)
    if err != nil {
        log.Fatal(err)
    }
    fmt.Println(osm)
}

Osm 和其他结构类似于您期望的 XML 模式。 decoder.Decode(&amp;osm) 应用该架构。

如果您只想提取部分 XML,see the answers to How to extract part of an XML file as a string?。

答案的其余部分将仅涉及通道和 goroutines 的使用。 XML 部分将被删除。


如果你添加一些调试语句,你会发现parseRelation 从未被调用,这意味着channel 是空的,fmt.Println(&lt;- channel) 坐在那里等待一个永远不会关闭的空通道。因此,一旦您完成处理,请关闭频道。

  for {
    p2Rc, err := reader.ReadSlice('\n')

    ...
  }
  close(channel)

现在我们得到{ [] []} 5973974 次。

for i := 0; i < 5973974; i++ {
  fmt.Println(<- channel)
}

这是尝试从频道读取 5973974 次。这违背了渠道的意义。相反,read from the channel using range。

for thing := range channel {
    fmt.Println(thing)
}

现在至少它完成了!

但是有一个新问题。如果它真的找到了一个东西,比如如果你将if p2Rc[2] == 114 { 更改为if p2Rc[2] == 32 {,你会得到一个panic: send on closed channel。这是因为parseRelation 与阅读器并行运行,可能会在主阅读代码完成并关闭通道后尝试写入。您必须确保在关闭频道之前使用该频道的每个人都已完成。

要解决这个问题,需要进行相当大的重新设计。


这是一个简单程序的示例,它从文件中读取行,将它们放入通道中,并让工作人员从该通道中读取数据。

func main() {
    reader, err := os.Open("test.xml")
    if err != nil {
        log.Fatal(err)
    }
    defer reader.Close()

    // Use the simpler bufio.Scanner
    scanner := bufio.NewScanner(reader)

    // A buffered channel for input
    lines := make(chan string, 10)

    // Work on the lines
    go func() {
        for line := range lines {
            fmt.Println(line)
        }
    }()

    // Read lines into the channel
    for scanner.Scan() {
        lines <- scanner.Text()
    }
    if err := scanner.Err(); err != nil {
        log.Fatal(err)
    }

    // When main exits, channels gracefully close.
}

这很好用,因为main 是特殊的,它在退出时会清理通道。但是如果 reader 和 writer 都是 goroutine 呢?

// A buffered channel for input
lines := make(chan string, 10)

// Work on the lines
go func() {
    for line := range lines {
        fmt.Println(line)
    }
}()

// Read lines into the channel
go func() {
    for scanner.Scan() {
        lines <- scanner.Text()
    }
    if err := scanner.Err(); err != nil {
        log.Fatal(err)
    }
}()

空的。 main 在 goroutine 完成工作之前退出并关闭通道。我们需要一种方法让main 知道等待处理完成。有几种方法可以做到这一点。一种是another channel to synchronize processing。

// A buffered channel for input
lines := make(chan string, 10)

// A channel for main to wait for
done := make(chan bool, 1)

// Work on the lines
go func() {
    for line := range lines {
        fmt.Println(line)
    }

    // Indicate the worker is done
    done <- true
}()

// Read lines into the channel
go func() {
    // Explicitly close `lines` when we're done so the workers can finish
    defer close(lines)

    for scanner.Scan() {
        lines <- scanner.Text()
    }
    if err := scanner.Err(); err != nil {
        log.Fatal(err)
    }
}()

// main will wait until there's something to read from `done`
<-done

现在main 将触发阅读器和工作 goroutines 和缓冲区,等待 done 上的某些内容。阅读器将填写lines,直到完成阅读,然后关闭它。并行工作人员将从lines 读取并在完成读取后写入done。

另一个选项是使用sync.WaitGroup。

// A buffered channel for input
lines := make(chan string, 10)

var wg sync.WaitGroup

// Note that there is one more thing to wait for
wg.Add(1)
go func() {
    // Tell the WaitGroup we're done
    defer wg.Done()

    for line := range lines {
        fmt.Println(line)
    }
}()

// Read lines into the channel
go func() {
    defer close(lines)

    for scanner.Scan() {
        lines <- scanner.Text()
    }
    if err := scanner.Err(); err != nil {
        log.Fatal(err)
    }
}()

// Wait until everything in the WaitGroup is done
wg.Wait()

和以前一样,main 启动 reader 和 worker goroutine,但现在它在启动 worker 之前将 1 添加到 WaitGroup。然后它一直等到wg.Wait() 返回。阅读器的工作方式与以前相同,完成后关闭lines 频道。工作人员现在在完成递减 WaitGroup 的计数并允许 wg.Wait() 返回时调用 wg.Done()。

每种技术都有优点和缺点。 done 更灵活,链更好,如果你能把头绕在它周围,它会更安全。 WaitGroups 更简单,更容易理解,但要求每个 goroutine 共享一个变量。


如果我们想添加到这个处理链中,我们可以这样做。假设我们有一个读取行的 goroutine,一个在 XML 元素中处理它们,一个对元素做一些事情。

// A buffered channel for input
lines := make(chan []byte, 10)
elements := make(chan *RS)

var wg sync.WaitGroup

// Worker goroutine, turn lines into RS structs
wg.Add(1)
go func() {
    defer wg.Done()
    defer close(elements)

    for line := range lines {
        if line[2] == 32 {
            fmt.Println("Begin")
            fmt.Println(string(line))
            fmt.Println("End")

            rs := &RS{}
            xml.Unmarshal(line, &rs)
            elements <- rs
        }
    }
}()

// File reader
go func() {
    defer close(lines)

    for scanner.Scan() {
        lines <- scanner.Bytes()
    }
    if err := scanner.Err(); err != nil {
        log.Fatal(err)
    }
}()

// Element reader
wg.Add(1)
go func() {
    defer wg.Done()

    for element := range elements {
        fmt.Println(element)
    }
}()

wg.Wait()

这会产生空结构,因为您试图将 XML 的各个行推入表示完整 &lt;relationship&gt; 标记的结构中。但它展示了如何向链中添加更多工人。

【讨论】:

  • 我想你的意思是在//Element reader 中添加defer wg.Done(),除了那个很棒的回复,谢谢!
  • 多么有趣的断言'首先,让我们解决这个问题:您不能逐行解析 XML。 ' 我猜堆栈的概念可能已经过时了
  • @Adonis Newline is not a statement terminator in XML,它只是另一个空白字符。每行可以有多个标签,标签和字符串可以跨越多行。
【解决方案2】:

我不知道这是否会帮助你,但你的条件,if p2Rc[2]==114 永远不会满足,那么你继续并开始收听频道。它永远不会收到输入。此外,还有更好的收听频道的方法,例如select,这是一个示例https://tour.golang.org/concurrency/5。

我认为这里的主要问题是 this(code) 的目的是什么?上述条件在哪里属于该过程?如果这很清楚,我可以更新一个更好的响应。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2021-06-26
    • 1970-01-01
    • 2019-07-30
    • 2013-08-13
    • 2021-12-10
    • 2021-01-29
    • 2015-05-20
    • 1970-01-01
    相关资源
    最近更新 更多